diff --git a/modules/disk-load/README.md b/modules/disk-load/README.md new file mode 100644 index 00000000..4ba20e5e --- /dev/null +++ b/modules/disk-load/README.md @@ -0,0 +1,87 @@ +# disk-load + +How loaded a machine's storage is, read from `/proc` and `/sys` and nothing else: the share of CPU time +spent waiting on I/O, I/O pressure stall, each disk's busy share, throughput, IOPS, wait and queue +depth, the filesystems' space, and the units, containers and processes doing the I/O (novox/hq issue +312). It changes nothing and declares no resources, so it can be assigned to every machine. + +## Why it exists, and why here + +Investigating why the controller's checks ran five to seven times slower on the control node (hq issue +306) needed disk numbers, and no tool in the mesh reported any: not iowait, not a disk's utilisation, +not throughput, latency or which process wrote. The operator's rule is that a missing tool is built in +the module that owns the area, never worked round with a login to the machine. + +No module owned the area. `memory-pressure` reads the same kind of thing for memory, but it owns +compressed swap and systemd-oomd: assigning it to a machine to read the disks would change how that +machine handles memory. No node seat holds the machine's storage either. Every seat is a role one +module holds on each machine, with verbs that act through it, and a new seat costs a controller change +first (hq ADR 0246). A report that reads kernel counters needs neither a role nor exclusivity. So this +is a small module with tools only, like `netcheck` beside the network. If the mesh later gives a +machine's resources a seat, these tools become its verbs. + +## Tools + +All four only read. + +| tool | what | replaces | +|---|---|---| +| `disk_load` | Samples a window (`seconds`, default 2, at most 10). It answers the CPU's iowait, user, system and idle %; I/O pressure (some and full, over the window and as avg10/60/300); per disk busy %, read and write MiB/s, IOPS, wait per request in ms, queue depth and requests in flight, with the filesystems each disk holds; the filesystems' space; and the busiest units and processes (`processes`, default 5, 0 for none). `said` puts each finding in one line, flagged past a line (below) | `iostat -x`, `vmstat`, `iotop`, `df` | +| `disk_top` | Over a window. `by=unit` (the default) ranks systemd units and containers by bytes moved, from their cgroups. `by=container` ranks containers by name, and `by=process` ranks processes | `iotop` | +| `disk_filesystems` | Each filesystem that stores data, said once however many places it is mounted: size, used and free GiB, used % as `df` counts it, inodes used % | `df -h`, `df -i` | +| `disk_devices` | Each disk without sampling: model, size, whether it spins, scheduler, queue limit, what it is built on, and since boot the GiB read and written, the busy share and the wait per read and per write | `lsblk`, `cat /proc/diskstats` | + +### How each number is read + +- **iowait** is from `/proc/stat`'s `cpu` line: the share of every CPU's time in the window spent idle + with I/O outstanding. Guest time is left out, because the kernel already counts it in user. +- **A disk** is read from `/proc/diskstats` at both ends of the window, with iostat's arithmetic. + Busy is the share of the window with any request in flight. Wait is the time spent on the + requests completed, divided by their number. Queue depth is the weighted time divided by the window. + Sectors there are always 512 bytes. A counter that went backwards (a device removed and added) + reads as 0, never as a huge number. Loop devices, RAM disks and zram are left out unless `all`; + zram is memory, and `memory-pressure` reports it. +- **Pressure** is `/proc/pressure/io`. The share of the window comes from the `total` counter, in + microseconds. A kernel without PSI is said as not available. +- **Units and containers** come from each cgroup's `io.stat`, which every account can read. Only the + innermost `.service` or `.scope` counts, because `io.stat` includes descendants. Only whole disks + built on nothing count: a write through an encrypted or logical volume is charged to the volume and + again to the disk under it, and a loop device's writes are charged again to the filesystem its file + sits on. Container names come from `docker ps`; if the runtime cannot be asked, a container is + given by its short id. +- **Processes** come from `/proc//io` (`read_bytes`, `write_bytes`, which is storage, not the + page cache). Another account's process can be read only when the tool runner is root. The unreadable + are counted and said, and the units above miss none of them. A process that started inside the + window is left out. +- **Filesystems** come from the process's own mount table, keeping only types that store data + (ext4, xfs, btrfs, vfat, NFS and the like). Each one is checked with statfs, given 2 s per mount so + a dead network mount costs one line. + +### The lines + +| finding | past | +|---|---| +| iowait HIGH | 20 % of CPU time | +| I/O pressure HIGH | some task stalled 10 % of the window, or avg60 ≥ 10 % | +| disk SATURATED | busy 90 % of the window | +| filesystem NEARLY FULL | 90 % of its space or of its inodes | + +Busy % is how iostat counts it, so it overstates an NVMe disk, which serves many requests at once. Read +it beside the wait and the queue depth. + +## Tests + +`go test ./...`. The tests read a tree standing in for `/proc` and `/sys`. Its fixture is a server +writing 100 MiB in 2 s to an encrypted root on an NVMe disk, with a spinning disk idle beside it. They +cover: + +- iowait, each disk's numbers, and the encrypted root named, with what it sits on and what it holds; +- loop, zram and partitions left out unless asked for; +- pressure, and a kernel without it; +- a counter reset, and busy held at 100; +- units counted once, by innermost cgroup and bottom disk; +- containers by name, or by short id if the runtime cannot be asked; +- processes ranked, one that started inside the window left out, unreadable ones counted; +- filesystems said once, mount escapes undone, a statfs that fails said; +- the window's bounds; +- the manifest: tools listed equal tools served, all read-only, no resources, no machine named. diff --git a/modules/disk-load/cmd/disk-load/cgroups.go b/modules/disk-load/cmd/disk-load/cgroups.go new file mode 100644 index 00000000..79af261b --- /dev/null +++ b/modules/disk-load/cmd/disk-load/cgroups.go @@ -0,0 +1,169 @@ +package main + +import ( + "io/fs" + "path/filepath" + "sort" + "strconv" + "strings" + "time" +) + +// cgroupIO is a cgroup's bytes and requests to storage, from its io.stat. Unlike /proc//io, io.stat +// is readable by every account, so a unit or a container is seen whoever runs the tool runner. +type cgroupIO struct { + read, write, rios, wios uint64 +} + +// parseIOStat sums the lines of an io.stat ("259:0 rbytes=… wbytes=… rios=… wios=…") over the devices +// counted counts; every device when counted is nil. +func parseIOStat(s string, counted map[string]bool) cgroupIO { + var c cgroupIO + for _, line := range strings.Split(s, "\n") { + f := strings.Fields(line) + if len(f) < 2 || (counted != nil && !counted[f[0]]) { + continue + } + for _, kv := range f[1:] { + k, v, _ := strings.Cut(kv, "=") + n, _ := strconv.ParseUint(v, 10, 64) + switch k { + case "rbytes": + c.read += n + case "wbytes": + c.write += n + case "rios": + c.rios += n + case "wios": + c.wios += n + } + } + } + return c +} + +// bottomDevices are the major:minor of the disks I/O finally lands on. A write through an encrypted +// or logical volume is charged to the volume and again to the disk under it, and a loop device's to the +// filesystem its file is on, so only a whole disk built on nothing is counted. nil (count every line) +// when /proc/diskstats cannot be read. +func (m *Machine) bottomDevices() map[string]bool { + text := m.read("/proc/diskstats") + if text == "" { + return nil + } + stats := parseDiskstats(text) + disks := map[string]bool{} + for _, name := range m.blockDevices(stats, false) { + if len(m.below(name)) == 0 { + disks[name] = true + } + } + out := map[string]bool{} + for _, line := range strings.Split(text, "\n") { + f := strings.Fields(line) + if len(f) > 2 && disks[f[2]] { + out[f[0]+":"+f[1]] = true + } + } + return out +} + +// unitIO is every unit's storage I/O, by its cgroup's path: the innermost .service or .scope, because +// io.stat counts a cgroup's descendants too and an outer unit would count them twice. +func (m *Machine) unitIO(counted map[string]bool) map[string]cgroupIO { + base := m.path("/sys/fs/cgroup") + var units []string + filepath.WalkDir(base, func(p string, d fs.DirEntry, err error) error { + if err != nil || !d.IsDir() { + return nil + } + if n := d.Name(); strings.HasSuffix(n, ".service") || strings.HasSuffix(n, ".scope") { + rel, _ := filepath.Rel(base, p) + units = append(units, filepath.ToSlash(rel)) + } + return nil + }) + outer := map[string]bool{} + for _, u := range units { + for dir := filepath.Dir(u); dir != "." && dir != "/"; dir = filepath.Dir(dir) { + outer[dir] = true + } + } + out := map[string]cgroupIO{} + for _, u := range units { + if outer[u] { + continue + } + text := m.read("/sys/fs/cgroup/" + u + "/io.stat") + if text == "" { + continue + } + out[u] = parseIOStat(text, counted) + } + return out +} + +// topUnits is the units, or the containers, that moved the most bytes between two readings of io.stat. +func (m *Machine) topUnits(a, b map[string]cgroupIO, window time.Duration, by string, limit int) []Process { + s := window.Seconds() + groups := map[string]*Process{} + var names map[string]string + for path, end := range b { + start, ok := a[path] + if !ok { + continue + } + r, w := sub(end.read, start.read), sub(end.write, start.write) + if r+w == 0 { + continue + } + unit := filepath.Base(path) + key := unit + var container string + if id := containerID.FindStringSubmatch("/" + path); id != nil { + if names == nil { + names = m.containerNames() + } + container = id[1][:12] + if n, ok := names[id[1]]; ok { + container = n + } + } + if by == "container" { + if container == "" { + continue + } + key = container + } + g, ok := groups[key] + if !ok { + g = &Process{Unit: unit, Container: container} + if by == "container" { + g.Unit = "" + } + groups[key] = g + } + g.ReadMiB += r + g.WriteMiB += w + g.bytes += r + w + g.ReadIOPS += round1(sub(end.rios, start.rios) / s) + g.WriteIOPS += round1(sub(end.wios, start.wios) / s) + } + out := make([]Process, 0, len(groups)) + for _, g := range groups { + out = append(out, *g) + } + sort.Slice(out, func(i, j int) bool { + if out[i].bytes != out[j].bytes { + return out[i].bytes > out[j].bytes + } + return out[i].Unit+out[i].Container < out[j].Unit+out[j].Container + }) + if len(out) > limit { + out = out[:limit] + } + for i := range out { + inMiB(&out[i], s) + } + return out +} diff --git a/modules/disk-load/cmd/disk-load/filesystems.go b/modules/disk-load/cmd/disk-load/filesystems.go new file mode 100644 index 00000000..b80e0ddc --- /dev/null +++ b/modules/disk-load/cmd/disk-load/filesystems.go @@ -0,0 +1,132 @@ +package main + +import ( + "path/filepath" + "strings" +) + +// storedOn are the filesystem types that hold data on a device or a server. Everything else in the +// mount table (proc, sysfs, cgroup, tmpfs, overlay layers, squashfs images…) is the kernel's or a copy. +var storedOn = map[string]bool{ + "ext2": true, "ext3": true, "ext4": true, "xfs": true, "btrfs": true, "zfs": true, "f2fs": true, + "bcachefs": true, "vfat": true, "exfat": true, "ntfs": true, "ntfs3": true, "fuseblk": true, + "nfs": true, "nfs4": true, "cifs": true, "smb3": true, +} + +// mount is one line of the mount table. +type mount struct{ source, target, fstype string } + +// unescape undoes the mount table's octal escapes (a space is \040). +func unescape(s string) string { + if !strings.Contains(s, `\`) { + return s + } + var b strings.Builder + for i := 0; i < len(s); i++ { + if s[i] == '\\' && i+3 < len(s) { + o := s[i+1 : i+4] + if len(o) == 3 && o[0] >= '0' && o[0] <= '3' && o[1] >= '0' && o[1] <= '7' && o[2] >= '0' && o[2] <= '7' { + b.WriteByte((o[0]-'0')<<6 | (o[1]-'0')<<3 | (o[2] - '0')) + i += 3 + continue + } + } + b.WriteByte(s[i]) + } + return b.String() +} + +// mounts is this process's mount table, only the filesystems that store data. +func (m *Machine) mounts() []mount { + var out []mount + for _, line := range strings.Split(m.read("/proc/self/mounts"), "\n") { + f := strings.Fields(line) + if len(f) < 3 || !storedOn[f[2]] { + continue + } + out = append(out, mount{source: unescape(f[0]), target: unescape(f[1]), fstype: f[2]}) + } + return out +} + +// Filesystem is one filesystem's space, once however many places it is mounted at. +type Filesystem struct { + Mount string `json:"mount"` + AlsoAt []string `json:"also_at,omitempty"` + Source string `json:"source"` + Device string `json:"device,omitempty"` + Type string `json:"type"` + SizeGiB float64 `json:"size_gib"` + UsedGiB float64 `json:"used_gib"` + AvailGiB float64 `json:"available_gib"` + UsedPct float64 `json:"used_pct"` + InodesUsed float64 `json:"inodes_used_pct,omitempty"` + Unanswered string `json:"error,omitempty"` +} + +// blockName is the kernel's name for the device a mount's source names: /dev/nvme0n1p3 is nvme0n1p3, +// /dev/mapper/root the dm-N whose name is root; "" for a source that is no local device. +func (m *Machine) blockName(source string, mappers map[string]string) string { + if !strings.HasPrefix(source, "/dev/") { + return "" + } + if name, ok := strings.CutPrefix(source, "/dev/mapper/"); ok { + return mappers[name] + } + // /dev/disk/by-uuid/… and the like are links the system made; reading where one points creates nothing. + if m.Root == "/" { + if real, err := filepath.EvalSymlinks(source); err == nil { + source = real + } + } + return filepath.Base(source) +} + +// mappers is each device-mapper name with the dm-N it is. +func (m *Machine) mappers() map[string]string { + out := map[string]string{} + for _, dev := range m.names("/sys/block") { + if name := m.mapperName(dev); name != "" { + out[name] = dev + } + } + return out +} + +// Filesystems is the space of every filesystem that stores data. A filesystem mounted at several +// places (a btrfs subvolume, a bind mount) is said once, at its first mount, with the others beside it. +func (m *Machine) Filesystems() []Filesystem { + mappers := m.mappers() + var out []Filesystem + seen := map[string]int{} + for _, mt := range m.mounts() { + key := mt.source + if i, ok := seen[key]; ok { + out[i].AlsoAt = append(out[i].AlsoAt, mt.target) + continue + } + fs := Filesystem{Mount: mt.target, Source: mt.source, Type: mt.fstype, Device: m.blockName(mt.source, mappers)} + st, err := m.Statfs(m.path(mt.target)) + if err != nil { + fs.Unanswered = err.Error() + } else if st.Size > 0 { + used := st.Size - st.Free + fs.SizeGiB = round1(float64(st.Size) / gib) + fs.UsedGiB = round1(float64(used) / gib) + fs.AvailGiB = round1(float64(st.Avail) / gib) + // As df counts it: used against what an ordinary user could have, so the reserve shows as full. + if used+st.Avail > 0 { + fs.UsedPct = round1(float64(used) / float64(used+st.Avail) * 100) + } + if st.Files > 0 { + fs.InodesUsed = round1(float64(st.Files-st.FilesFree) / float64(st.Files) * 100) + } + } + seen[key] = len(out) + out = append(out, fs) + } + if out == nil { + out = []Filesystem{} + } + return out +} diff --git a/modules/disk-load/cmd/disk-load/helpers_test.go b/modules/disk-load/cmd/disk-load/helpers_test.go new file mode 100644 index 00000000..71a65655 --- /dev/null +++ b/modules/disk-load/cmd/disk-load/helpers_test.go @@ -0,0 +1,84 @@ +package main + +import ( + "context" + "fmt" + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +// fake is a machine for a test: a tree standing in for /proc and /sys, a statfs answering from a table, +// a runner answering from a table, and a wait that turns the tree into its second reading. +type fake struct { + t *testing.T + root string + later map[string]string + gone []string + space map[string]FSStat + answers map[string]string + calls []string + waited time.Duration +} + +func newFake(t *testing.T) *fake { + t.Helper() + return &fake{t: t, root: t.TempDir(), later: map[string]string{}, space: map[string]FSStat{}, answers: map[string]string{}} +} + +func (f *fake) machine() *Machine { + return &Machine{ + Root: f.root, + Run: func(_ context.Context, name string, args ...string) (string, error) { + line := strings.TrimSpace(name + " " + strings.Join(args, " ")) + f.calls = append(f.calls, line) + if out, ok := f.answers[line]; ok { + return out, nil + } + return "", fmt.Errorf("%s: not found", name) + }, + Wait: func(d time.Duration) { + f.waited += d + for p, c := range f.later { + f.file(p, c) + } + for _, p := range f.gone { + os.RemoveAll(filepath.Join(f.root, p)) + } + }, + Statfs: func(path string) (FSStat, error) { + rel, _ := filepath.Rel(f.root, path) + s, ok := f.space[filepath.Clean("/"+filepath.ToSlash(rel))] + if !ok { + return FSStat{}, fmt.Errorf("no such filesystem") + } + return s, nil + }, + } +} + +// file writes a file under the fake root. +func (f *fake) file(path, content string) { + f.t.Helper() + full := filepath.Join(f.root, path) + if err := os.MkdirAll(filepath.Dir(full), 0o755); err != nil { + f.t.Fatal(err) + } + if err := os.WriteFile(full, []byte(content), 0o644); err != nil { + f.t.Fatal(err) + } +} + +// dir makes a directory under the fake root: sysfs's entries that are links on a real machine are +// directories here, because the module only lists their names. +func (f *fake) dir(path string) { + f.t.Helper() + if err := os.MkdirAll(filepath.Join(f.root, path), 0o755); err != nil { + f.t.Fatal(err) + } +} + +// then is a file's content at the second reading. +func (f *fake) then(path, content string) { f.later[path] = content } diff --git a/modules/disk-load/cmd/disk-load/load.go b/modules/disk-load/cmd/disk-load/load.go new file mode 100644 index 00000000..61ae0b3c --- /dev/null +++ b/modules/disk-load/cmd/disk-load/load.go @@ -0,0 +1,264 @@ +package main + +import ( + "fmt" + "sort" + "strconv" + "strings" + "time" +) + +// The lines past which a reading is said as a finding, not only a number. +const ( + busyLine = 90.0 // a device with a request in flight this share of the window is saturated + iowaitLine = 20.0 // the processors idle on I/O this share of their time + stallLine = 10.0 // some task stalled on I/O this share of the window + fullLine = 90.0 // a filesystem this full, or this short of inodes +) + +// reading is everything read at one moment of the window. +type reading struct { + cpu cpuTimes + disks map[string]diskCounters + psi PSI + psiOK bool + procs map[int]procIO + cantRead int + units map[string]cgroupIO +} + +func (m *Machine) readNow(withProcesses bool, counted map[string]bool) (reading, error) { + var r reading + var err error + if r.cpu, err = parseCPU(m.read("/proc/stat")); err != nil { + return r, err + } + r.disks = parseDiskstats(m.read("/proc/diskstats")) + r.psi, r.psiOK = parsePSI(m.read("/proc/pressure/io")) + if withProcesses { + r.procs, r.cantRead = m.processIO() + r.units = m.unitIO(counted) + } + return r, nil +} + +// Load is the machine's disk load over one window. +type Load struct { + WindowS float64 `json:"window_s"` + Said []string `json:"said"` + CPU CPU `json:"cpu"` + Pressure any `json:"io_pressure"` + Devices []Device `json:"devices"` + Filesystems []Filesystem `json:"filesystems,omitempty"` + Units []Process `json:"units,omitempty"` + Processes []Process `json:"processes,omitempty"` + Unreadable int `json:"processes_unreadable,omitempty"` + Note string `json:"note,omitempty"` +} + +// LoadOver reads the machine, waits the window, reads it again and says what happened between. +// processes is how many of the busiest processes to name (0 for none); all keeps every device and +// partition, not only the whole disks. +func (m *Machine) LoadOver(window time.Duration, processes int, all, filesystems bool) (Load, error) { + counted := m.bottomDevices() + a, err := m.readNow(processes > 0, counted) + if err != nil { + return Load{}, err + } + m.Wait(window) + b, err := m.readNow(processes > 0, counted) + if err != nil { + return Load{}, err + } + l := Load{WindowS: window.Seconds(), CPU: cpuShare(a.cpu, b.cpu)} + + mounts := map[string][]string{} + if filesystems { + l.Filesystems = m.Filesystems() + for _, fs := range l.Filesystems { + if fs.Device == "" { + continue + } + mounts[fs.Device] = append(mounts[fs.Device], fs.Mount) + if p := m.parentOf(fs.Device); p != fs.Device { + mounts[p] = append(mounts[p], fs.Mount+" ("+fs.Device+")") + } + } + } + l.Devices = []Device{} + for _, name := range m.blockDevices(b.disks, all) { + start, ok := a.disks[name] + if !ok { + continue + } + d := deviceOver(name, start, b.disks[name], window) + d.Name = m.mapperName(name) + d.On = m.below(name) + d.Mounts = mounts[name] + l.Devices = append(l.Devices, d) + } + sort.SliceStable(l.Devices, func(i, j int) bool { return l.Devices[i].BusyPct > l.Devices[j].BusyPct }) + + if a.psiOK && b.psiOK { + l.Pressure = psiOver(a.psi, b.psi, window) + } else { + l.Pressure = "not available: this kernel has no /proc/pressure/io (pressure stall information is off)" + } + if processes > 0 { + l.Units = m.topUnits(a.units, b.units, window, "unit", processes) + l.Processes = m.topProcesses(a.procs, b.procs, window, processes) + l.Unreadable = b.cantRead + if b.cantRead > 0 { + l.Note = unreadableNote(b.cantRead) + } + } + l.Said = say(l) + return l, nil +} + +func say(l Load) []string { + var out []string + c := l.CPU + line := fmt.Sprintf("CPU: %.1f %% waiting on I/O, %.1f %% user, %.1f %% system, %.1f %% idle, over %g s on %d CPUs", c.IOWaitPct, c.UserPct, c.SystemPct, c.IdlePct, l.WindowS, c.CPUs) + if c.IOWaitPct >= iowaitLine { + line += " — HIGH: the processors sat idle waiting on storage" + } + out = append(out, line) + if p, ok := l.Pressure.(PSI); ok { + line := fmt.Sprintf("I/O pressure: some task stalled %.1f %% of the window (avg10 %.1f %%, avg60 %.1f %%, avg300 %.1f %%)", p.Some.WindowPct, p.Some.Avg10, p.Some.Avg60, p.Some.Avg300) + if p.Full != nil { + line += fmt.Sprintf("; every task stalled %.1f %%", p.Full.WindowPct) + } + if p.Some.WindowPct >= stallLine || p.Some.Avg60 >= stallLine { + line += " — HIGH" + } + out = append(out, line) + } + for _, d := range l.Devices { + name := d.Device + if d.Name != "" { + name += " (" + d.Name + ")" + } + line := fmt.Sprintf("%s: busy %.1f %%, read %.1f MiB/s (%.1f IOPS, %.1f ms wait), write %.1f MiB/s (%.1f IOPS, %.1f ms wait), queue depth %.1f", + name, d.BusyPct, d.ReadMiBs, d.ReadIOPS, d.ReadWaitMs, d.WriteMiBs, d.WriteIOPS, d.WriteWaitMs, d.QueueDepth) + if len(d.Mounts) > 0 { + line += ", holds " + strings.Join(d.Mounts, ", ") + } + if d.BusyPct >= busyLine { + line += " — SATURATED" + } + out = append(out, line) + } + for _, fs := range l.Filesystems { + if fs.Unanswered != "" { + out = append(out, fmt.Sprintf("%s: not read (%s)", fs.Mount, fs.Unanswered)) + continue + } + if fs.UsedPct >= fullLine || fs.InodesUsed >= fullLine { + out = append(out, fmt.Sprintf("%s: %.1f %% full (%.1f GiB of %.1f GiB, %.1f GiB free), inodes %.1f %% used — NEARLY FULL", fs.Mount, fs.UsedPct, fs.UsedGiB, fs.SizeGiB, fs.AvailGiB, fs.InodesUsed)) + } + } + for i, u := range l.Units { + if i == 3 { + break + } + name := u.Unit + if u.Container != "" { + name = "container " + u.Container + } + out = append(out, fmt.Sprintf("busiest unit: %s read %.2f MiB/s, wrote %.2f MiB/s", name, u.ReadMiBs, u.WriteMiBs)) + } + for i, p := range l.Processes { + if i == 3 { + break + } + where := p.Unit + if p.Container != "" { + where = "container " + p.Container + } + out = append(out, fmt.Sprintf("busiest process: %s (pid %d, %s) read %.2f MiB/s, wrote %.2f MiB/s", p.Command, p.PID, orNone(where), p.ReadMiBs, p.WriteMiBs)) + } + return out +} + +// unreadableNote says what a list of processes misses when the tool runner is not root. +func unreadableNote(n int) string { + return fmt.Sprintf("%d processes' I/O could not be read: they are another account's, and the tool runner is not root, "+ + "so a busy one may be missing from the processes; the units, read from their cgroups, miss none", n) +} + +func orNone(s string) string { + if s == "" { + return "no unit" + } + return s +} + +// DeviceInfo is one disk as it is, and what it has done since the machine started. +type DeviceInfo struct { + Device string `json:"device"` + Name string `json:"name,omitempty"` + On []string `json:"on,omitempty"` + Model string `json:"model,omitempty"` + SizeGiB float64 `json:"size_gib"` + Rotational bool `json:"rotational"` + Scheduler string `json:"scheduler,omitempty"` + QueueLimit int `json:"queue_limit,omitempty"` + ReadGiB float64 `json:"read_since_boot_gib"` + WrittenGiB float64 `json:"written_since_boot_gib"` + BusySincePct float64 `json:"busy_since_boot_pct"` + ReadWaitMs float64 `json:"read_wait_since_boot_ms"` + WriteWaitMs float64 `json:"write_wait_since_boot_ms"` +} + +// DevicesNow is every disk with its kind, its size, its scheduler and its totals since boot: the +// question "is this a spinning disk" and "how much has it written" without a window. +func (m *Machine) DevicesNow(all bool) ([]DeviceInfo, float64) { + stats := parseDiskstats(m.read("/proc/diskstats")) + uptime := 0.0 + if f := strings.Fields(m.read("/proc/uptime")); len(f) > 0 { + uptime, _ = strconv.ParseFloat(f[0], 64) + } + out := []DeviceInfo{} + for _, name := range m.blockDevices(stats, all) { + c := stats[name] + base := "/sys/block/" + name + if !m.exists(base) { + base = "/sys/block/" + m.parentOf(name) + "/" + name + } + sectors, _ := strconv.ParseFloat(strings.TrimSpace(m.read(base+"/size")), 64) + d := DeviceInfo{ + Device: name, + Name: m.mapperName(name), + On: m.below(name), + Model: strings.TrimSpace(m.read(base + "/device/model")), + SizeGiB: round1(sectors * 512 / gib), + Rotational: strings.TrimSpace(m.read(base+"/queue/rotational")) == "1", + Scheduler: activeScheduler(m.read(base + "/queue/scheduler")), + ReadGiB: round1(float64(c.readSectors) * 512 / gib), + WrittenGiB: round1(float64(c.writeSectors) * 512 / gib), + } + d.QueueLimit, _ = strconv.Atoi(strings.TrimSpace(m.read(base + "/queue/nr_requests"))) + if uptime > 0 { + d.BusySincePct = round1(float64(c.busyMs) / (uptime * 1000) * 100) + } + if c.reads > 0 { + d.ReadWaitMs = round1(float64(c.readMs) / float64(c.reads)) + } + if c.writes > 0 { + d.WriteWaitMs = round1(float64(c.writeMs) / float64(c.writes)) + } + out = append(out, d) + } + return out, round1(uptime / 3600) +} + +// activeScheduler is the bracketed one of "mq-deadline kyber [bfq] none". +func activeScheduler(s string) string { + if i := strings.Index(s, "["); i >= 0 { + if j := strings.Index(s[i:], "]"); j > 0 { + return s[i+1 : i+j] + } + } + return strings.TrimSpace(s) +} diff --git a/modules/disk-load/cmd/disk-load/load_test.go b/modules/disk-load/cmd/disk-load/load_test.go new file mode 100644 index 00000000..7927f231 --- /dev/null +++ b/modules/disk-load/cmd/disk-load/load_test.go @@ -0,0 +1,307 @@ +package main + +import ( + "encoding/json" + "os" + "sort" + "strings" + "testing" + "time" +) + +const store = "4f1c2a9e8b7d6c5a4f1c2a9e8b7d6c5a4f1c2a9e8b7d6c5a4f1c2a9e8b7d6c5a" + +// aServerUnderWrites is a machine whose store writes 100 MiB in a 2 s window to an encrypted root on an +// NVMe disk: the shape of the control node in issue 312. Lines are as a 6.x kernel writes them, with +// the four discard and two flush fields at the end of /proc/diskstats. +func aServerUnderWrites(t *testing.T) *fake { + f := newFake(t) + f.file("/proc/stat", "cpu 1000 0 500 8000 100 0 0 0 0 0\ncpu0 500 0 250 4000 50 0 0 0 0 0\ncpu1 500 0 250 4000 50 0 0 0 0 0\nintr 1\n") + f.then("/proc/stat", "cpu 1100 0 550 8250 500 0 0 0 0 0\ncpu0 550 0 275 4125 250 0 0 0 0 0\ncpu1 550 0 275 4125 250 0 0 0 0 0\nintr 2\n") + + f.file("/proc/diskstats", strings.Join([]string{ + " 259 0 nvme0n1 1000 0 80000 500 2000 0 400000 4000 0 10000 20000 0 0 0 0 10 5", + " 259 1 nvme0n1p1 10 0 100 5 0 0 0 0 0 5 5 0 0 0 0 0 0", + " 259 2 nvme0n1p2 990 0 79900 495 2000 0 400000 4000 0 9995 19995 0 0 0 0 0 0", + " 254 0 dm-0 990 0 79900 600 2000 0 400000 5000 0 10000 21000 0 0 0 0 0 0", + " 8 0 sda 50 0 4000 900 0 0 0 0 0 800 900 0 0 0 0 0 0", + " 7 0 loop0 5 0 10 1 0 0 0 0 0 1 1 0 0 0 0 0 0", + " 252 0 zram0 100 0 800 1 900 0 7200 9 0 9 10 0 0 0 0 0 0", + }, "\n")+"\n") + f.then("/proc/diskstats", strings.Join([]string{ + " 259 0 nvme0n1 1200 0 120960 900 3000 0 604800 29000 12 11960 45400 0 0 0 0 12 6", + " 259 1 nvme0n1p1 10 0 100 5 0 0 0 0 0 5 5 0 0 0 0 0 0", + " 259 2 nvme0n1p2 1190 0 120860 895 3000 0 604800 29000 12 11955 45395 0 0 0 0 0 0", + " 254 0 dm-0 1190 0 120860 1100 3000 0 604800 32000 14 11980 50000 0 0 0 0 0 0", + " 8 0 sda 50 0 4000 900 0 0 0 0 0 800 900 0 0 0 0 0 0", + " 7 0 loop0 5 0 10 1 0 0 0 0 0 1 1 0 0 0 0 0 0", + " 252 0 zram0 900 0 7200 9 9000 0 72000 90 0 90 100 0 0 0 0 0 0", + }, "\n")+"\n") + + f.file("/proc/pressure/io", "some avg10=38.50 avg60=12.00 avg300=3.10 total=1000000\nfull avg10=25.00 avg60=8.00 avg300=2.00 total=500000\n") + f.then("/proc/pressure/io", "some avg10=44.00 avg60=14.20 avg300=3.50 total=1900000\nfull avg10=29.00 avg60=9.10 avg300=2.20 total=1100000\n") + + f.file("/sys/block/nvme0n1/size", "1000215216\n") + f.file("/sys/block/nvme0n1/queue/rotational", "0\n") + f.file("/sys/block/nvme0n1/queue/scheduler", "[none] mq-deadline kyber\n") + f.file("/sys/block/nvme0n1/queue/nr_requests", "1023\n") + f.file("/sys/block/nvme0n1/device/model", "Example NVMe 512GB \n") + f.dir("/sys/block/nvme0n1/nvme0n1p1") + f.dir("/sys/block/nvme0n1/nvme0n1p2") + f.file("/sys/block/dm-0/dm/name", "root\n") + f.dir("/sys/block/dm-0/slaves/nvme0n1p2") + f.file("/sys/block/dm-0/size", "999000000\n") + f.file("/sys/block/sda/size", "7814037168\n") + f.file("/sys/block/sda/queue/rotational", "1\n") + f.file("/sys/block/sda/queue/scheduler", "mq-deadline [bfq] none\n") + f.dir("/sys/block/sda/sda1") + f.dir("/sys/block/loop0") + f.dir("/sys/block/zram0") + f.file("/proc/uptime", "36000.00 70000.00\n") + + f.file("/proc/self/mounts", strings.Join([]string{ + "proc /proc proc rw,nosuid 0 0", + "/dev/mapper/root / ext4 rw,relatime 0 0", + "tmpfs /tmp tmpfs rw 0 0", + "/dev/nvme0n1p1 /boot vfat rw 0 0", + "/dev/mapper/root /var/lib/docker ext4 rw,relatime 0 0", + "overlay /var/lib/docker/overlay2/abc/merged overlay rw 0 0", + `/dev/sda1 /mnt/back\040up ext4 rw 0 0`, + }, "\n")+"\n") + f.space["/"] = FSStat{Size: 100 * gib, Free: 10 * gib, Avail: 5 * gib, Files: 1000, FilesFree: 600} + f.space["/boot"] = FSStat{Size: gib, Free: gib / 2, Avail: gib / 2, Files: 0} + + proc := func(pid, comm, cgroup string, r0, w0, r1, w1 int) { + f.file("/proc/"+pid+"/comm", comm+"\n") + f.file("/proc/"+pid+"/cgroup", cgroup+"\n") + f.file("/proc/"+pid+"/io", ioFile(r0, w0)) + f.then("/proc/"+pid+"/io", ioFile(r1, w1)) + } + proc("100", "postgres", "0::/system.slice/docker-"+store+".scope", 0, 0, 0, 100*mib) + proc("101", "postgres", "0::/system.slice/docker-"+store+".scope", 0, 0, 10*mib, 0) + proc("200", "systemd-journal", "0::/system.slice/systemd-journald.service", 0, 0, 0, 2*mib) + proc("201", "sleep", "0::/user.slice/user-1000.slice/session-1.scope", 0, 0, 0, 0) + f.dir("/proc/300") // another account's: its io is not readable + f.then("/proc/400/io", ioFile(0, 50*mib)) + cg := func(path, at0, at1 string) { + f.file("/sys/fs/cgroup/"+path+"/io.stat", at0) + f.then("/sys/fs/cgroup/"+path+"/io.stat", at1) + } + zero := "259:0 rbytes=0 wbytes=0 rios=0 wios=0 dbytes=0 dios=0\n254:0 rbytes=0 wbytes=0 rios=0 wios=0 dbytes=0 dios=0\n7:0 rbytes=0 wbytes=0 rios=0 wios=0 dbytes=0 dios=0\n" + // The store's writes are charged to the encrypted volume and again to the disk under it; a loop + // device's to the filesystem its file is on. Only the disk is counted. + cg("system.slice/docker-"+store+".scope", zero, + "259:0 rbytes="+itoa(10*mib)+" wbytes="+itoa(100*mib)+" rios=200 wios=1000 dbytes=0 dios=0\n"+ + "254:0 rbytes="+itoa(10*mib)+" wbytes="+itoa(100*mib)+" rios=200 wios=1000 dbytes=0 dios=0\n"+ + "7:0 rbytes=0 wbytes="+itoa(900*mib)+" rios=0 wios=9000 dbytes=0 dios=0\n") + cg("system.slice/systemd-journald.service", zero, "259:0 rbytes=0 wbytes="+itoa(2*mib)+" rios=0 wios=20 dbytes=0 dios=0\n") + cg("system.slice/idle.service", zero, zero) + // io.stat counts a cgroup's descendants: the user's manager holds the application's bytes again. + cg("user.slice/user-1000.slice/user@1000.service", zero, "259:0 rbytes=0 wbytes="+itoa(5*mib)+" rios=0 wios=50 dbytes=0 dios=0\n") + cg("user.slice/user-1000.slice/user@1000.service/app.slice/editor.scope", zero, "259:0 rbytes=0 wbytes="+itoa(5*mib)+" rios=0 wios=50 dbytes=0 dios=0\n") + f.answers["docker ps --no-trunc --format {{.ID}}\t{{.Names}}"] = store + "\tstore-check\n" + return f +} + +func ioFile(read, write int) string { + return "rchar: 99999999\nwchar: 99999999\nsyscr: 1\nsyscw: 1\nread_bytes: " + itoa(read) + "\nwrite_bytes: " + itoa(write) + "\ncancelled_write_bytes: 0\n" +} + +func itoa(n int) string { + b, _ := json.Marshal(n) + return string(b) +} + +func device(t *testing.T, l Load, name string) Device { + t.Helper() + for _, d := range l.Devices { + if d.Device == name { + return d + } + } + t.Fatalf("no device %s in %+v", name, l.Devices) + return Device{} +} + +func TestAWindowSaysIOWaitEachDisksLoadAndPressure(t *testing.T) { + f := aServerUnderWrites(t) + l, err := f.machine().LoadOver(2*time.Second, 5, false, true) + if err != nil { + t.Fatal(err) + } + if f.waited != 2*time.Second { + t.Errorf("waited %s, not the window", f.waited) + } + if l.CPU != (CPU{CPUs: 2, IOWaitPct: 50, UserPct: 12.5, SystemPct: 6.3, IdlePct: 31.3}) { + t.Errorf("cpu %+v", l.CPU) + } + + n := device(t, l, "nvme0n1") + want := Device{Device: "nvme0n1", Mounts: []string{"/boot (nvme0n1p1)"}, BusyPct: 98, ReadMiBs: 10, WriteMiBs: 50, + ReadIOPS: 100, WriteIOPS: 500, ReadWaitMs: 2, WriteWaitMs: 25, QueueDepth: 12.7, InFlight: 12} + if got, _ := json.Marshal(n); string(got) != string(must(json.Marshal(want))) { + t.Errorf("nvme0n1\n got %s\nwant %s", got, must(json.Marshal(want))) + } + root := device(t, l, "dm-0") + if root.Name != "root" || strings.Join(root.On, ",") != "nvme0n1p2" || strings.Join(root.Mounts, ",") != "/" { + t.Errorf("the encrypted root is said by its name, what it is on and what it holds: %+v", root) + } + if sda := device(t, l, "sda"); sda.BusyPct != 0 || sda.ReadMiBs != 0 { + t.Errorf("an idle disk: %+v", sda) + } + for _, d := range l.Devices { + if d.Device == "loop0" || d.Device == "zram0" || d.Device == "nvme0n1p2" { + t.Errorf("%s is not a disk and was said", d.Device) + } + } + if l.Devices[0].Device != "dm-0" && l.Devices[0].Device != "nvme0n1" { + t.Errorf("the busiest first: %s", l.Devices[0].Device) + } + + p, ok := l.Pressure.(PSI) + if !ok || p.Some.WindowPct != 45 || p.Full == nil || p.Full.WindowPct != 30 || p.Some.Avg10 != 44 { + t.Errorf("pressure %+v", l.Pressure) + } + + if len(l.Processes) != 3 { + t.Fatalf("three processes moved bytes in the window, one started inside it: %+v", l.Processes) + } + top := l.Processes[0] + if top.PID != 100 || top.Command != "postgres" || top.Container != "store-check" || top.WriteMiBs != 50 || top.WriteMiB != 100 { + t.Errorf("the busiest process, named by its container: %+v", top) + } + if j := l.Processes[2]; j.Unit != "systemd-journald.service" || j.Container != "" || j.WriteMiBs != 1 { + t.Errorf("a service's process, named by its unit: %+v", j) + } + if len(l.Units) != 3 || l.Units[0].WriteIOPS != 500 || l.Units[2].Unit != "systemd-journald.service" { + t.Errorf("units %+v", l.Units) + } + if l.Unreadable != 1 || !strings.Contains(l.Note, "1 processes' I/O could not be read") { + t.Errorf("an unreadable process is counted and said: %d %q", l.Unreadable, l.Note) + } + + said := strings.Join(l.Said, "\n") + for _, s := range []string{ + "CPU: 50.0 % waiting on I/O", "HIGH: the processors sat idle waiting on storage", + "I/O pressure: some task stalled 45.0 % of the window", "every task stalled 30.0 %", + "nvme0n1: busy 98.0 %, read 10.0 MiB/s (100.0 IOPS, 2.0 ms wait), write 50.0 MiB/s (500.0 IOPS, 25.0 ms wait), queue depth 12.7, holds /boot (nvme0n1p1) — SATURATED", + "dm-0 (root): busy 99.0 %", "holds / — SATURATED", + "/: 94.7 % full (90.0 GiB of 100.0 GiB, 5.0 GiB free)", "NEARLY FULL", + `/mnt/back up: not read (no such filesystem)`, + "busiest process: postgres (pid 100, container store-check) read 0.00 MiB/s, wrote 50.00 MiB/s", + "busiest unit: container store-check read 5.00 MiB/s, wrote 50.00 MiB/s", + "busiest unit: editor.scope read 0.00 MiB/s, wrote 2.50 MiB/s", + } { + if !strings.Contains(said, s) { + t.Errorf("said lacks %q:\n%s", s, said) + } + } +} + +func must(b []byte, err error) []byte { + if err != nil { + panic(err) + } + return b +} + +func TestAllSaysEveryDeviceAndPartition(t *testing.T) { + l, err := aServerUnderWrites(t).machine().LoadOver(2*time.Second, 0, true, false) + if err != nil { + t.Fatal(err) + } + var names []string + for _, d := range l.Devices { + names = append(names, d.Device) + } + sort.Strings(names) + if strings.Join(names, ",") != "dm-0,loop0,nvme0n1,nvme0n1p1,nvme0n1p2,sda,zram0" { + t.Errorf("all: %v", names) + } + if l.Processes != nil || l.Filesystems != nil { + t.Errorf("processes=0 names none and reads none: %+v %+v", l.Processes, l.Filesystems) + } +} + +func TestAKernelWithoutPressureStallSaysSo(t *testing.T) { + f := aServerUnderWrites(t) + os.Remove(f.root + "/proc/pressure/io") + delete(f.later, "/proc/pressure/io") + l, err := f.machine().LoadOver(time.Second, 0, false, false) + if err != nil { + t.Fatal(err) + } + if s, ok := l.Pressure.(string); !ok || !strings.HasPrefix(s, "not available") { + t.Errorf("pressure %+v", l.Pressure) + } +} + +func TestACounterThatWentBackwardsReadsAsNothingNotAsAnEnormousNumber(t *testing.T) { + a := diskCounters{reads: 100, readSectors: 1000, busyMs: 5000} + b := diskCounters{reads: 10, readSectors: 100, busyMs: 10} + d := deviceOver("sdx", a, b, time.Second) + if d.ReadIOPS != 0 || d.ReadMiBs != 0 || d.BusyPct != 0 { + t.Errorf("a reset counter: %+v", d) + } + if c := cpuShare(cpuTimes{iowait: 10, idle: 10}, cpuTimes{iowait: 5, idle: 10}); c.IOWaitPct != 0 { + t.Errorf("cpu over no time: %+v", c) + } +} + +func TestBusyIsHeldAtAHundred(t *testing.T) { + d := deviceOver("sdx", diskCounters{}, diskCounters{busyMs: 1500}, time.Second) + if d.BusyPct != 100 { + t.Errorf("busy %v", d.BusyPct) + } +} + +func TestAStatWithoutACPULineIsAnError(t *testing.T) { + if _, err := parseCPU("intr 1\n"); err == nil { + t.Error("no cpu line was accepted") + } + if _, err := parseCPU("cpu 1 2 3\n"); err == nil { + t.Error("a short cpu line was accepted") + } + f := newFake(t) + if _, err := f.machine().LoadOver(time.Second, 0, false, false); err == nil { + t.Error("a machine without /proc/stat answered") + } +} + +func TestDevicesSayWhatEachDiskIsAndHasDoneSinceBoot(t *testing.T) { + ds, up := aServerUnderWrites(t).machine().DevicesNow(false) + if up != 10 { + t.Errorf("up %v hours", up) + } + by := map[string]DeviceInfo{} + for _, d := range ds { + by[d.Device] = d + } + n := by["nvme0n1"] + if n.Model != "Example NVMe 512GB" || n.Rotational || n.Scheduler != "none" || n.QueueLimit != 1023 || n.SizeGiB != 476.9 { + t.Errorf("nvme0n1 %+v", n) + } + if n.WrittenGiB != 0.2 || n.ReadWaitMs != 0.5 || n.WriteWaitMs != 2 || n.BusySincePct != 0 { + t.Errorf("nvme0n1 since boot %+v", n) + } + if s := by["sda"]; !s.Rotational || s.Scheduler != "bfq" || s.SizeGiB != 3726 || s.BusySincePct != 0 { + t.Errorf("sda %+v", s) + } + if r := by["dm-0"]; r.Name != "root" || strings.Join(r.On, ",") != "nvme0n1p2" { + t.Errorf("dm-0 %+v", r) + } + if _, ok := by["zram0"]; ok { + t.Error("zram is memory, not a disk") + } +} + +func TestAPartitionsSizeIsReadUnderItsDisk(t *testing.T) { + f := aServerUnderWrites(t) + f.file("/sys/block/nvme0n1/nvme0n1p1/size", "2097152\n") + ds, _ := f.machine().DevicesNow(true) + for _, d := range ds { + if d.Device == "nvme0n1p1" && d.SizeGiB != 1 { + t.Errorf("nvme0n1p1 %+v", d) + } + } +} diff --git a/modules/disk-load/cmd/disk-load/machine.go b/modules/disk-load/cmd/disk-load/machine.go new file mode 100644 index 00000000..99c79c33 --- /dev/null +++ b/modules/disk-load/cmd/disk-load/machine.go @@ -0,0 +1,124 @@ +package main + +import ( + "bytes" + "context" + "fmt" + "os" + "os/exec" + "path/filepath" + "strings" + "syscall" + "time" +) + +// CommandTimeout bounds the one command this module runs (the container runtime's list of names). +const CommandTimeout = 5 * time.Second + +// StatfsTimeout bounds the look at one mount: a network filesystem whose server is gone blocks statfs +// for as long as it likes, and must cost the answer one line, never the whole call. +const StatfsTimeout = 2 * time.Second + +// Runner runs one command and answers its standard output. It is injected so that the tools are tested +// against recorded answers rather than this machine's daemons. +type Runner func(ctx context.Context, name string, args ...string) (string, error) + +// ExecRunner runs a command on the machine, bounded by CommandTimeout. +func ExecRunner(ctx context.Context, name string, args ...string) (string, error) { + ctx, cancel := context.WithTimeout(ctx, CommandTimeout) + defer cancel() + cmd := exec.CommandContext(ctx, name, args...) + var stdout, stderr bytes.Buffer + cmd.Stdout, cmd.Stderr = &stdout, &stderr + if err := cmd.Run(); err != nil { + said := strings.TrimSpace(stderr.String()) + if said != "" { + return "", fmt.Errorf("%s: %w: %s", name, err, said) + } + return "", fmt.Errorf("%s: %w", name, err) + } + return stdout.String(), nil +} + +// FSStat is what statfs says of one mounted filesystem, in bytes and inodes. +type FSStat struct { + Size, Free, Avail uint64 + Files, FilesFree uint64 +} + +// Machine is what the module reads: a filesystem root (the real one, or a test's tree standing in for +// /proc and /sys), a way to run a command, a way to wait out the sampling window, and statfs. +type Machine struct { + Root string + Run Runner + Wait func(time.Duration) + Statfs func(path string) (FSStat, error) +} + +// Here is the machine this process runs on. +func Here() *Machine { + return &Machine{Root: "/", Run: ExecRunner, Wait: time.Sleep, Statfs: statfs} +} + +func statfs(path string) (FSStat, error) { + type answer struct { + s FSStat + err error + } + done := make(chan answer, 1) + go func() { + var s syscall.Statfs_t + if err := syscall.Statfs(path, &s); err != nil { + done <- answer{err: err} + return + } + bs := uint64(s.Bsize) + done <- answer{s: FSStat{Size: s.Blocks * bs, Free: s.Bfree * bs, Avail: s.Bavail * bs, Files: s.Files, FilesFree: s.Ffree}} + }() + select { + case a := <-done: + return a.s, a.err + case <-time.After(StatfsTimeout): + return FSStat{}, fmt.Errorf("statfs did not answer within %s", StatfsTimeout) + } +} + +func (m *Machine) path(p string) string { return filepath.Join(m.Root, p) } + +// read is a file's content; "" when it cannot be read. +func (m *Machine) read(p string) string { + b, err := os.ReadFile(m.path(p)) + if err != nil { + return "" + } + return string(b) +} + +func (m *Machine) exists(p string) bool { + _, err := os.Stat(m.path(p)) + return err == nil +} + +// names lists the entries of a directory, sorted; none when it cannot be read. +func (m *Machine) names(dir string) []string { + entries, err := os.ReadDir(m.path(dir)) + if err != nil { + return nil + } + out := make([]string, 0, len(entries)) + for _, e := range entries { + out = append(out, e.Name()) + } + return out +} + +// round1 rounds to one decimal, for what a person reads. +func round1(f float64) float64 { + if f < 0 { + return -float64(int64(-f*10+0.5)) / 10 + } + return float64(int64(f*10+0.5)) / 10 +} + +const mib = 1 << 20 +const gib = 1 << 30 diff --git a/modules/disk-load/cmd/disk-load/main.go b/modules/disk-load/cmd/disk-load/main.go new file mode 100644 index 00000000..92fd2ef4 --- /dev/null +++ b/modules/disk-load/cmd/disk-load/main.go @@ -0,0 +1,20 @@ +// The disk-load module's Go bundle (novox/hq ADR 0188, ADR 0193, issue 312): a process the node's +// runtime launches and speaks MCP over stdio to. It reads the machine's storage load from /proc and +// /sys — iowait, pressure stall, each disk's busy share, throughput, IOPS, wait and queue depth, the +// filesystems' space and the processes doing the I/O — and changes nothing. +package main + +import ( + "fmt" + "os" + + stdio "git.novox.be/novox/mesh-sdk/go" +) + +func main() { + // An empty name serves as the module the runtime names (MESH_SERVED_MODULE): disk-load. + if err := stdio.Serve("", Tools(Here())); err != nil { + fmt.Fprintln(os.Stderr, err) + os.Exit(1) + } +} diff --git a/modules/disk-load/cmd/disk-load/processes.go b/modules/disk-load/cmd/disk-load/processes.go new file mode 100644 index 00000000..1ddfc58c --- /dev/null +++ b/modules/disk-load/cmd/disk-load/processes.go @@ -0,0 +1,173 @@ +package main + +import ( + "context" + "regexp" + "sort" + "strconv" + "strings" + "time" +) + +// procIO is what /proc//io says a process caused to be read from and written to storage: not +// what it read from the page cache, which rchar and wchar count. +type procIO struct { + read, write uint64 +} + +// processIO reads every process's storage counters. A process of another account is unreadable unless +// this runtime runs as root; those are counted, never guessed. +func (m *Machine) processIO() (map[int]procIO, int) { + out := map[int]procIO{} + unreadable := 0 + for _, name := range m.names("/proc") { + pid, err := strconv.Atoi(name) + if err != nil { + continue + } + text := m.read("/proc/" + name + "/io") + if text == "" { + if m.exists("/proc/" + name) { + unreadable++ + } + continue + } + var p procIO + for _, line := range strings.Split(text, "\n") { + k, v, ok := strings.Cut(line, ":") + if !ok { + continue + } + n, _ := strconv.ParseUint(strings.TrimSpace(v), 10, 64) + switch k { + case "read_bytes": + p.read = n + case "write_bytes": + p.write = n + } + } + out[pid] = p + } + return out, unreadable +} + +// Process is one process's storage I/O over the window, or a unit's or a container's (with its requests +// per second, which io.stat counts and /proc//io does not). +type Process struct { + PID int `json:"pid,omitempty"` + Command string `json:"command,omitempty"` + Unit string `json:"unit,omitempty"` + Container string `json:"container,omitempty"` + ReadMiBs float64 `json:"read_mib_s"` + WriteMiBs float64 `json:"write_mib_s"` + ReadMiB float64 `json:"read_mib"` + WriteMiB float64 `json:"write_mib"` + ReadIOPS float64 `json:"read_iops,omitempty"` + WriteIOPS float64 `json:"write_iops,omitempty"` + bytes float64 +} + +var containerID = regexp.MustCompile(`(?:docker-|/docker/)([0-9a-f]{64})`) + +// placeOf is the systemd unit and the container a process runs in, from its cgroup. +func (m *Machine) placeOf(pid int) (unit, container string) { + for _, line := range strings.Split(m.read("/proc/"+strconv.Itoa(pid)+"/cgroup"), "\n") { + parts := strings.SplitN(line, ":", 3) + if len(parts) != 3 || (parts[0] != "0" && !strings.Contains(parts[1], "name=systemd")) { + continue + } + path := strings.TrimSpace(parts[2]) + if id := containerID.FindStringSubmatch(path); id != nil { + container = id[1] + } + segs := strings.Split(path, "/") + for i := len(segs) - 1; i >= 0; i-- { + if strings.HasSuffix(segs[i], ".service") || strings.HasSuffix(segs[i], ".scope") { + unit = segs[i] + break + } + } + if parts[0] == "0" { + break + } + } + return unit, container +} + +// containerNames is each running container's full id with its name, from the container runtime; none +// when it cannot be asked, and then a container is said by its short id. +func (m *Machine) containerNames() map[string]string { + out := map[string]string{} + if m.Run == nil { + return out + } + text, err := m.Run(context.Background(), "docker", "ps", "--no-trunc", "--format", "{{.ID}}\t{{.Names}}") + if err != nil { + return out + } + for _, line := range strings.Split(text, "\n") { + id, name, ok := strings.Cut(strings.TrimSpace(line), "\t") + if ok { + out[id] = name + } + } + return out +} + +// topProcesses is the processes that moved the most bytes to and from storage between the two +// readings, the largest first, each with the unit and the container it runs in. +func (m *Machine) topProcesses(a, b map[int]procIO, window time.Duration, limit int) []Process { + var moved []Process + for pid, end := range b { + start, ok := a[pid] + if !ok { + continue + } + r, w := sub(end.read, start.read), sub(end.write, start.write) + if r+w == 0 { + continue + } + // Bytes until inMiB says them in MiB. + moved = append(moved, Process{PID: pid, ReadMiB: r, WriteMiB: w, bytes: r + w}) + } + sort.Slice(moved, func(i, j int) bool { + if moved[i].bytes != moved[j].bytes { + return moved[i].bytes > moved[j].bytes + } + return moved[i].PID < moved[j].PID + }) + if len(moved) > limit { + moved = moved[:limit] + } + var names map[string]string + for i := range moved { + p := &moved[i] + p.Command = strings.TrimSpace(m.read("/proc/" + strconv.Itoa(p.PID) + "/comm")) + var id string + p.Unit, id = m.placeOf(p.PID) + if id != "" { + if names == nil { + names = m.containerNames() + } + p.Container = id[:12] + if n, ok := names[id]; ok { + p.Container = n + } + } + inMiB(p, window.Seconds()) + } + if moved == nil { + moved = []Process{} + } + return moved +} + +// inMiB turns a reading's bytes into MiB and MiB/s, to two decimals: a process writing a log moves +// kilobytes, and one decimal of a MiB would say it moved nothing. +func inMiB(p *Process, s float64) { + r, w := p.ReadMiB, p.WriteMiB + p.ReadMiBs, p.WriteMiBs = round2(r/mib/s), round2(w/mib/s) + p.ReadMiB, p.WriteMiB = round2(r/mib), round2(w/mib) +} + +func round2(f float64) float64 { return float64(int64(f*100+0.5)) / 100 } diff --git a/modules/disk-load/cmd/disk-load/sample.go b/modules/disk-load/cmd/disk-load/sample.go new file mode 100644 index 00000000..2d46eb5d --- /dev/null +++ b/modules/disk-load/cmd/disk-load/sample.go @@ -0,0 +1,281 @@ +package main + +import ( + "fmt" + "sort" + "strconv" + "strings" + "time" +) + +// cpuTimes is the first line of /proc/stat, in clock ticks: what every CPU together spent its time on. +type cpuTimes struct { + user, nice, system, idle, iowait, irq, softirq, steal uint64 + cpus int +} + +// total leaves guest time out: the kernel already counts it in user and nice. +func (c cpuTimes) total() uint64 { + return c.user + c.nice + c.system + c.idle + c.iowait + c.irq + c.softirq + c.steal +} + +func parseCPU(stat string) (cpuTimes, error) { + var c cpuTimes + found := false + for _, line := range strings.Split(stat, "\n") { + f := strings.Fields(line) + if len(f) == 0 || !strings.HasPrefix(f[0], "cpu") { + continue + } + if f[0] != "cpu" { + c.cpus++ + continue + } + if len(f) < 9 { + return c, fmt.Errorf("/proc/stat: the cpu line has %d fields, not at least 9", len(f)) + } + v := make([]uint64, 8) + for i := range v { + n, err := strconv.ParseUint(f[i+1], 10, 64) + if err != nil { + return c, fmt.Errorf("/proc/stat: %q is not a count", f[i+1]) + } + v[i] = n + } + c.user, c.nice, c.system, c.idle, c.iowait, c.irq, c.softirq, c.steal = v[0], v[1], v[2], v[3], v[4], v[5], v[6], v[7] + found = true + } + if !found { + return c, fmt.Errorf("/proc/stat has no cpu line") + } + return c, nil +} + +// CPU is how all the processors together spent the window, in percent of their time. +type CPU struct { + CPUs int `json:"cpus"` + IOWaitPct float64 `json:"iowait_pct"` + UserPct float64 `json:"user_pct"` + SystemPct float64 `json:"system_pct"` + IdlePct float64 `json:"idle_pct"` + StealPct float64 `json:"steal_pct"` +} + +func cpuShare(a, b cpuTimes) CPU { + total := float64(b.total() - a.total()) + pct := func(x, y uint64) float64 { + if total <= 0 || y < x { + return 0 + } + return round1(float64(y-x) / total * 100) + } + return CPU{ + CPUs: b.cpus, + IOWaitPct: pct(a.iowait, b.iowait), + UserPct: pct(a.user+a.nice, b.user+b.nice), + SystemPct: pct(a.system+a.irq+a.softirq, b.system+b.irq+b.softirq), + IdlePct: pct(a.idle, b.idle), + StealPct: pct(a.steal, b.steal), + } +} + +// diskCounters is one line of /proc/diskstats (Documentation/admin-guide/iostats.rst). Sectors are +// always 512 bytes there, whatever the device's own sector size. +type diskCounters struct { + reads, readSectors, readMs uint64 + writes, writeSectors, writeMs uint64 + inFlight, busyMs, weightedMs uint64 +} + +func parseDiskstats(s string) map[string]diskCounters { + out := map[string]diskCounters{} + for _, line := range strings.Split(s, "\n") { + f := strings.Fields(line) + if len(f) < 14 { + continue + } + n := make([]uint64, 11) + ok := true + for i := range n { + v, err := strconv.ParseUint(f[i+3], 10, 64) + if err != nil { + ok = false + break + } + n[i] = v + } + if !ok { + continue + } + out[f[2]] = diskCounters{ + reads: n[0], readSectors: n[2], readMs: n[3], + writes: n[4], writeSectors: n[6], writeMs: n[7], + inFlight: n[8], busyMs: n[9], weightedMs: n[10], + } + } + return out +} + +// Device is one block device over the window. +type Device struct { + Device string `json:"device"` + Name string `json:"name,omitempty"` + On []string `json:"on,omitempty"` + Mounts []string `json:"mounts,omitempty"` + BusyPct float64 `json:"busy_pct"` + ReadMiBs float64 `json:"read_mib_s"` + WriteMiBs float64 `json:"write_mib_s"` + ReadIOPS float64 `json:"read_iops"` + WriteIOPS float64 `json:"write_iops"` + ReadWaitMs float64 `json:"read_wait_ms"` + WriteWaitMs float64 `json:"write_wait_ms"` + QueueDepth float64 `json:"queue_depth"` + InFlight uint64 `json:"in_flight"` +} + +func sub(b, a uint64) float64 { + if b < a { + return 0 + } + return float64(b - a) +} + +// deviceOver is iostat's arithmetic: busy is the share of the window with any request in flight, wait +// the time a completed request took from queueing to completion, queue depth the requests in flight +// on average. +func deviceOver(name string, a, b diskCounters, window time.Duration) Device { + ms := float64(window.Milliseconds()) + s := window.Seconds() + reads, writes := sub(b.reads, a.reads), sub(b.writes, a.writes) + wait := func(spent, n float64) float64 { + if n == 0 { + return 0 + } + return round1(spent / n) + } + busy := sub(b.busyMs, a.busyMs) / ms * 100 + if busy > 100 { + busy = 100 + } + return Device{ + Device: name, + BusyPct: round1(busy), + ReadMiBs: round1(sub(b.readSectors, a.readSectors) * 512 / mib / s), + WriteMiBs: round1(sub(b.writeSectors, a.writeSectors) * 512 / mib / s), + ReadIOPS: round1(reads / s), + WriteIOPS: round1(writes / s), + ReadWaitMs: wait(sub(b.readMs, a.readMs), reads), + WriteWaitMs: wait(sub(b.writeMs, a.writeMs), writes), + QueueDepth: round1(sub(b.weightedMs, a.weightedMs) / ms), + InFlight: b.inFlight, + } +} + +// blockDevices are the whole devices in /sys/block, without the ones that are not a disk: loop devices, +// RAM disks and compressed swap in RAM (zram, the memory-pressure module's). all keeps every one, and +// every partition /proc/diskstats names. +func (m *Machine) blockDevices(stats map[string]diskCounters, all bool) []string { + var out []string + if all { + for name := range stats { + out = append(out, name) + } + sort.Strings(out) + return out + } + for _, name := range m.names("/sys/block") { + if strings.HasPrefix(name, "loop") || strings.HasPrefix(name, "ram") || strings.HasPrefix(name, "zram") { + continue + } + if _, ok := stats[name]; ok { + out = append(out, name) + } + } + return out +} + +// mapperName is a device-mapper device's own name (dm-0 is "root"), "" for any other device. +func (m *Machine) mapperName(dev string) string { + return strings.TrimSpace(m.read("/sys/block/" + dev + "/dm/name")) +} + +// below is what a device-mapper or RAID device is built on. +func (m *Machine) below(dev string) []string { return m.names("/sys/block/" + dev + "/slaves") } + +// parentOf is the whole device a partition is on; the device itself when it is whole. +func (m *Machine) parentOf(dev string) string { + if m.exists("/sys/block/" + dev) { + return dev + } + for _, disk := range m.names("/sys/block") { + if m.exists("/sys/block/" + disk + "/" + dev) { + return disk + } + } + return dev +} + +// PSI is the kernel's pressure stall information for I/O: the share of time some (or all) runnable +// tasks were stalled waiting on it. +type PSI struct { + Some PSILine `json:"some"` + Full *PSILine `json:"full,omitempty"` +} + +// PSILine is one line, its running averages and its share of the sampled window. +type PSILine struct { + WindowPct float64 `json:"window_pct"` + Avg10 float64 `json:"avg10_pct"` + Avg60 float64 `json:"avg60_pct"` + Avg300 float64 `json:"avg300_pct"` + totalUs uint64 +} + +func parsePSI(s string) (PSI, bool) { + var p PSI + some := false + for _, line := range strings.Split(s, "\n") { + f := strings.Fields(line) + if len(f) < 5 { + continue + } + var l PSILine + for _, kv := range f[1:] { + k, v, _ := strings.Cut(kv, "=") + switch k { + case "avg10": + l.Avg10, _ = strconv.ParseFloat(v, 64) + case "avg60": + l.Avg60, _ = strconv.ParseFloat(v, 64) + case "avg300": + l.Avg300, _ = strconv.ParseFloat(v, 64) + case "total": + l.totalUs, _ = strconv.ParseUint(v, 10, 64) + } + } + switch f[0] { + case "some": + p.Some, some = l, true + case "full": + full := l + p.Full = &full + } + } + return p, some +} + +func psiOver(a, b PSI, window time.Duration) PSI { + share := func(x, y PSILine) PSILine { + y.WindowPct = round1(sub(y.totalUs, x.totalUs) / float64(window.Microseconds()) * 100) + if y.WindowPct > 100 { + y.WindowPct = 100 + } + return y + } + out := PSI{Some: share(a.Some, b.Some)} + if a.Full != nil && b.Full != nil { + full := share(*a.Full, *b.Full) + out.Full = &full + } + return out +} diff --git a/modules/disk-load/cmd/disk-load/tools.go b/modules/disk-load/cmd/disk-load/tools.go new file mode 100644 index 00000000..b54bf21e --- /dev/null +++ b/modules/disk-load/cmd/disk-load/tools.go @@ -0,0 +1,172 @@ +package main + +import ( + "fmt" + "math" + "strconv" + "strings" + "time" + + stdio "git.novox.be/novox/mesh-sdk/go" +) + +// Tools is the module's tools over one machine. Every one only reads. +func Tools(m *Machine) []stdio.Tool { + window := map[string]any{"type": "integer", "description": "how long to sample, in seconds (default 2, at most 10)"} + return []stdio.Tool{ + { + Name: "disk_load", + Description: "How loaded this machine's storage is, sampled over a few seconds: the share of CPU time spent waiting on I/O (iowait), " + + "I/O pressure stall (PSI), and for each disk how busy it was (%), read and write throughput (MiB/s), IOPS, the average wait " + + "per request (ms) and the queue depth, with the filesystems it holds; then the filesystems' space and the processes " + + "and the units moving the most bytes (from their cgroups), with their container. Each finding is said in a line under `said`. Replaces iostat, vmstat, " + + "iotop and df. (r)", + Input: map[string]any{ + "seconds": window, + "processes": map[string]any{"type": "integer", "description": "how many of the busiest units and processes to name (default 5, 0 for none, at most 50)"}, + "all": map[string]any{"type": "boolean", "description": "every block device and partition, not only the whole disks (default false)"}, + }, + Run: func(args map[string]any) (any, error) { + w, err := seconds(args) + if err != nil { + return nil, err + } + n, err := count(args, "processes", 5, 0, 50) + if err != nil { + return nil, err + } + return m.LoadOver(w, n, flag(args, "all"), true) + }, + }, + { + Name: "disk_top", + Description: "What moved the most bytes to and from storage over a few seconds: by=unit (the default) each systemd unit " + + "and container from its cgroup's io.stat, with MiB/s and IOPS, complete whoever runs the tools; by=container only the " + + "containers, by name; by=process each process from /proc//io, which misses another account's unless the tool " + + "runner is root. Replaces iotop. (r)", + Input: map[string]any{ + "seconds": window, + "limit": map[string]any{"type": "integer", "description": "how many (default 10, at most 50)"}, + "by": map[string]any{"type": "string", "enum": []string{"process", "unit", "container"}, "description": "what to rank (default unit)"}, + }, + Run: func(args map[string]any) (any, error) { + w, err := seconds(args) + if err != nil { + return nil, err + } + limit, err := count(args, "limit", 10, 1, 50) + if err != nil { + return nil, err + } + by := text(args, "by") + switch by { + case "": + by = "unit" + case "process", "unit", "container": + default: + return nil, fmt.Errorf("by is process, unit or container, not %q", by) + } + out := map[string]any{"window_s": w.Seconds(), "by": by} + if by != "process" { + counted := m.bottomDevices() + a := m.unitIO(counted) + m.Wait(w) + out["top"] = m.topUnits(a, m.unitIO(counted), w, by, limit) + return out, nil + } + a, _ := m.processIO() + m.Wait(w) + b, unreadable := m.processIO() + out["top"] = m.topProcesses(a, b, w, limit) + if unreadable > 0 { + out["processes_unreadable"] = unreadable + out["note"] = unreadableNote(unreadable) + } + return out, nil + }, + }, + { + Name: "disk_filesystems", + Description: "The space of every filesystem that stores data (not proc, tmpfs or overlay layers): size, used and free in GiB, " + + "used % as df counts it, and inodes used %. A filesystem mounted at several places is said once. Replaces df. (r)", + Run: func(map[string]any) (any, error) { + fss := m.Filesystems() + said := []string{} + for _, fs := range fss { + switch { + case fs.Unanswered != "": + said = append(said, fmt.Sprintf("%s: not read (%s)", fs.Mount, fs.Unanswered)) + default: + line := fmt.Sprintf("%s (%s, %s): %.1f %% used, %.1f GiB free of %.1f GiB", fs.Mount, fs.Type, fs.Source, fs.UsedPct, fs.AvailGiB, fs.SizeGiB) + if fs.UsedPct >= fullLine || fs.InodesUsed >= fullLine { + line += " — NEARLY FULL" + } + said = append(said, line) + } + } + return map[string]any{"said": said, "filesystems": fss}, nil + }, + }, + { + Name: "disk_devices", + Description: "Each disk as it is, without sampling: model, size (GiB), whether it spins, its I/O scheduler and queue limit, what " + + "it is built on (for LVM, encryption, RAID), and since the machine started: GiB read and written, the share of the time it " + + "was busy and the average wait per read and per write (ms). Replaces lsblk and cat /proc/diskstats. (r)", + Input: map[string]any{ + "all": map[string]any{"type": "boolean", "description": "every block device and partition, not only the whole disks (default false)"}, + }, + Run: func(args map[string]any) (any, error) { + ds, up := m.DevicesNow(flag(args, "all")) + return map[string]any{"up_hours": up, "devices": ds}, nil + }, + }, + } +} + +func seconds(args map[string]any) (time.Duration, error) { + n, err := count(args, "seconds", 2, 1, 10) + return time.Duration(n) * time.Second, err +} + +func text(args map[string]any, key string) string { + s, _ := args[key].(string) + return strings.TrimSpace(s) +} + +func flag(args map[string]any, key string) bool { + switch v := args[key].(type) { + case bool: + return v + case string: + return v == "true" + } + return false +} + +// count is a whole-number argument, defaulted, refused below least, held at most. +func count(args map[string]any, key string, fallback, least, most int) (int, error) { + v, given := args[key] + if !given || v == nil { + return fallback, nil + } + var n int + switch x := v.(type) { + case float64: + if x != math.Trunc(x) { + return 0, fmt.Errorf("%s must be a whole number, not %v", key, x) + } + n = int(x) + case string: + i, err := strconv.Atoi(strings.TrimSpace(x)) + if err != nil { + return 0, fmt.Errorf("%s must be a whole number, not %q", key, x) + } + n = i + default: + return 0, fmt.Errorf("%s must be a whole number", key) + } + if n < least { + return 0, fmt.Errorf("%s must be at least %d", key, least) + } + return min(n, most), nil +} diff --git a/modules/disk-load/cmd/disk-load/tools_test.go b/modules/disk-load/cmd/disk-load/tools_test.go new file mode 100644 index 00000000..88147715 --- /dev/null +++ b/modules/disk-load/cmd/disk-load/tools_test.go @@ -0,0 +1,200 @@ +package main + +import ( + "encoding/json" + "os" + "sort" + "strings" + "testing" + "time" + + stdio "git.novox.be/novox/mesh-sdk/go" +) + +func tool(t *testing.T, m *Machine, name string) stdio.Tool { + t.Helper() + for _, tl := range Tools(m) { + if tl.Name == name { + return tl + } + } + t.Fatalf("no tool %s", name) + return stdio.Tool{} +} + +func TestTheManifestNamesExactlyTheToolsTheBundleServesAndNoMachine(t *testing.T) { + raw, err := os.ReadFile("../../module.json") + if err != nil { + t.Fatal(err) + } + var m struct { + Tools []string `json:"tools"` + Resources []any `json:"resources"` + } + if err := json.Unmarshal(raw, &m); err != nil { + t.Fatal(err) + } + var served []string + for _, tl := range Tools(newFake(t).machine()) { + served = append(served, tl.Name) + if !strings.HasSuffix(tl.Description, "(r)") { + t.Errorf("%s only reads and does not say so", tl.Name) + } + } + sort.Strings(served) + sort.Strings(m.Tools) + if strings.Join(served, ",") != strings.Join(m.Tools, ",") { + t.Fatalf("served %v, listed %v", served, m.Tools) + } + if len(m.Resources) != 0 { + t.Errorf("the module only reads, and declares %d resources", len(m.Resources)) + } + for _, banned := range []string{"g14", "novox", "ace", "shanks", "jochen", "/home/"} { + if strings.Contains(string(raw), banned) { + t.Errorf("the manifest says %q", banned) + } + } +} + +func TestDiskLoadSamplesTwoSecondsAndNamesFiveProcessesByDefault(t *testing.T) { + f := aServerUnderWrites(t) + out, err := tool(t, f.machine(), "disk_load").Run(map[string]any{}) + if err != nil { + t.Fatal(err) + } + l := out.(Load) + if f.waited != 2*time.Second || l.WindowS != 2 || len(l.Processes) != 3 || len(l.Units) != 3 || len(l.Filesystems) != 3 { + t.Errorf("waited %s, %+v", f.waited, l) + } +} + +func TestTheWindowIsAWholeNumberOfSecondsHeldAtTen(t *testing.T) { + f := aServerUnderWrites(t) + load := tool(t, f.machine(), "disk_load") + for _, bad := range []any{0.0, -1.0, 1.5, "soon", true} { + if _, err := load.Run(map[string]any{"seconds": bad}); err == nil { + t.Errorf("seconds %v was accepted", bad) + } + } + if _, err := load.Run(map[string]any{"seconds": 30.0}); err != nil || f.waited != 10*time.Second { + t.Errorf("30 s was not held at 10: waited %s, %v", f.waited, err) + } + if _, err := load.Run(map[string]any{"processes": -1.0}); err == nil { + t.Error("processes -1 was accepted") + } +} + +func TestDiskTopRanksUnitsFromTheirCgroupsCountingEachByteOnce(t *testing.T) { + f := aServerUnderWrites(t) + out, err := tool(t, f.machine(), "disk_top").Run(map[string]any{}) + if err != nil { + t.Fatal(err) + } + if out.(map[string]any)["by"] != "unit" || f.waited != 2*time.Second { + t.Errorf("the default: %+v, waited %s", out, f.waited) + } + got := out.(map[string]any)["top"].([]Process) + want := []Process{ + {Unit: "docker-" + store + ".scope", Container: "store-check", ReadMiBs: 5, WriteMiBs: 50, ReadMiB: 10, WriteMiB: 100, ReadIOPS: 100, WriteIOPS: 500}, + {Unit: "editor.scope", WriteMiBs: 2.5, WriteMiB: 5, WriteIOPS: 25}, + {Unit: "systemd-journald.service", WriteMiBs: 1, WriteMiB: 2, WriteIOPS: 10}, + } + if g, w := must(json.Marshal(got)), must(json.Marshal(want)); string(g) != string(w) { + t.Errorf("by unit\n got %s\nwant %s", g, w) + } +} + +func TestDiskTopByContainerNamesOnlyContainers(t *testing.T) { + out, err := tool(t, aServerUnderWrites(t).machine(), "disk_top").Run(map[string]any{"by": "container", "seconds": 1.0}) + if err != nil { + t.Fatal(err) + } + got := out.(map[string]any)["top"].([]Process) + if len(got) != 1 || got[0].Container != "store-check" || got[0].Unit != "" || got[0].WriteMiBs != 100 { + t.Errorf("by container %+v", got) + } +} + +func TestDiskTopByProcessSaysWhatItCouldNotRead(t *testing.T) { + f := aServerUnderWrites(t) + top := tool(t, f.machine(), "disk_top") + out, err := top.Run(map[string]any{"by": "process", "limit": 2.0}) + if err != nil { + t.Fatal(err) + } + got := out.(map[string]any)["top"].([]Process) + if len(got) != 2 || got[0].PID != 100 || got[1].PID != 101 || got[1].ReadMiBs != 5 || got[1].Command != "postgres" { + t.Errorf("by process %+v", got) + } + if out.(map[string]any)["processes_unreadable"] != 1 { + t.Errorf("unreadable not said: %+v", out) + } + if _, err := top.Run(map[string]any{"by": "disk"}); err == nil { + t.Error("by=disk was accepted") + } +} + +func TestAContainerTheRuntimeCannotNameIsSaidByItsShortID(t *testing.T) { + f := aServerUnderWrites(t) + delete(f.answers, "docker ps --no-trunc --format {{.ID}}\t{{.Names}}") + out, err := tool(t, f.machine(), "disk_top").Run(map[string]any{"by": "process"}) + if err != nil { + t.Fatal(err) + } + if got := out.(map[string]any)["top"].([]Process); got[0].Container != store[:12] { + t.Errorf("%+v", got[0]) + } + if n := len(f.calls); n != 1 { + t.Errorf("the runtime was asked %d times, not once", n) + } +} + +func TestFilesystemsAreSaidOnceWithTheirOtherMounts(t *testing.T) { + out, err := tool(t, aServerUnderWrites(t).machine(), "disk_filesystems").Run(nil) + if err != nil { + t.Fatal(err) + } + fss := out.(map[string]any)["filesystems"].([]Filesystem) + var mounts []string + for _, fs := range fss { + mounts = append(mounts, fs.Mount) + } + if strings.Join(mounts, ",") != "/,/boot,/mnt/back up" { + t.Fatalf("mounts %v: proc, tmpfs and overlay are not stored data, and a second mount is not a second filesystem", mounts) + } + root := fss[0] + if strings.Join(root.AlsoAt, ",") != "/var/lib/docker" || root.Device != "dm-0" || root.Type != "ext4" || + root.SizeGiB != 100 || root.UsedGiB != 90 || root.AvailGiB != 5 || root.UsedPct != 94.7 || root.InodesUsed != 40 { + t.Errorf("root %+v", root) + } + if fss[1].Device != "nvme0n1p1" || fss[1].UsedPct != 50 || fss[1].InodesUsed != 0 { + t.Errorf("boot %+v", fss[1]) + } + if fss[2].Unanswered == "" { + t.Errorf("a filesystem statfs could not read says so: %+v", fss[2]) + } + said := strings.Join(out.(map[string]any)["said"].([]string), "\n") + if !strings.Contains(said, "/ (ext4, /dev/mapper/root): 94.7 % used, 5.0 GiB free of 100.0 GiB — NEARLY FULL") || + !strings.Contains(said, "/boot (vfat, /dev/nvme0n1p1): 50.0 % used, 0.5 GiB free of 1.0 GiB\n") { + t.Errorf("said:\n%s", said) + } +} + +func TestTheMountTablesEscapesAreUndone(t *testing.T) { + for in, want := range map[string]string{`/a\040b`: "/a b", `/tab\011x`: "/tab\tx", `/plain`: "/plain", `/end\04`: `/end\04`, `/back\134slash`: `/back\slash`} { + if got := unescape(in); got != want { + t.Errorf("%q: %q, not %q", in, got, want) + } + } +} + +func TestDiskDevicesAnswersWithoutWaiting(t *testing.T) { + f := aServerUnderWrites(t) + out, err := tool(t, f.machine(), "disk_devices").Run(map[string]any{}) + if err != nil { + t.Fatal(err) + } + if f.waited != 0 || len(out.(map[string]any)["devices"].([]DeviceInfo)) != 3 { + t.Errorf("waited %s, %+v", f.waited, out) + } +} diff --git a/modules/disk-load/go.mod b/modules/disk-load/go.mod new file mode 100644 index 00000000..edf8078e --- /dev/null +++ b/modules/disk-load/go.mod @@ -0,0 +1,5 @@ +module diskload + +go 1.22 + +require git.novox.be/novox/mesh-sdk/go v0.1.6 diff --git a/modules/disk-load/go.sum b/modules/disk-load/go.sum new file mode 100644 index 00000000..0dd60617 --- /dev/null +++ b/modules/disk-load/go.sum @@ -0,0 +1,2 @@ +git.novox.be/novox/mesh-sdk/go v0.1.6 h1:9qzdYONYbJdWcu6sxQcq9v1LI0JxcfkiKYkMUzJSkVQ= +git.novox.be/novox/mesh-sdk/go v0.1.6/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY= diff --git a/modules/disk-load/module.json b/modules/disk-load/module.json new file mode 100644 index 00000000..3b7ce911 --- /dev/null +++ b/modules/disk-load/module.json @@ -0,0 +1,25 @@ +{ + "module": "disk-load", + "version": "1", + "tools": [ + "disk_load", + "disk_top", + "disk_filesystems", + "disk_devices" + ], + "build": { + "artifacts": [ + { + "name": "tools-go", + "kind": "bundle", + "language": "go", + "system": "arch", + "from": "cmd/disk-load", + "binary": "disk-load", + "loads": [ + "disk-load" + ] + } + ] + } +}