From ad4ba21a4f416e441af762648234b9b5d377efa6 Mon Sep 17 00:00:00 2001 From: nabil1440 <52530910+nabil1440@users.noreply.github.com> Date: Thu, 24 Sep 2026 11:44:23 +0600 Subject: [PATCH 1/2] feat(agent): send the disk activity, and conform to contract v0.4.0 Contract v0.4.0, "Disk activity": the bytes and the operations that the hardware disks read and wrote in the minute, from /proc/diskstats, and their peaks each second within the minute. null in the first sample, after a reboot, after a counter went back, and without a hardware disk. The README pins contract v0.4.0. Closes #43 --- README.md | 4 +- internal/agent/clean.go | 2 + internal/agent/config.go | 2 +- internal/agent/wire/wire.go | 13 +- internal/metrics/disk.go | 98 +++++++++++ internal/metrics/disk_test.go | 309 ++++++++++++++++++++++++++++++++++ internal/metrics/metrics.go | 18 +- internal/metrics/parse.go | 37 ++++ internal/metrics/windows.go | 9 +- 9 files changed, 477 insertions(+), 15 deletions(-) create mode 100644 internal/metrics/disk.go create mode 100644 internal/metrics/disk_test.go diff --git a/README.md b/README.md index d055ac7..bfd7f11 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.3.1. +Conforms to the FlyWP monitoring agent contract v0.4.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, and it sends the values and the server status (restart needed, waiting updates, OS, kernel, uptime) 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. +`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) 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. 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 a1dceae..782c39b 100644 --- a/internal/agent/clean.go +++ b/internal/agent/clean.go @@ -40,6 +40,8 @@ func cleanSample(s wire.Sample) wire.Sample { } for _, v := range []**uint64{ &s.MemoryUsedMaxBytes, &s.SwapUsedMaxBytes, &s.NetInMaxBytesPerSecond, &s.NetOutMaxBytesPerSecond, + &s.DiskReadBytes, &s.DiskWriteBytes, &s.DiskReadOps, &s.DiskWriteOps, + &s.DiskReadMaxBytesPerSecond, &s.DiskWriteMaxBytesPerSecond, &s.DiskReadMaxOpsPerSecond, &s.DiskWriteMaxOpsPerSecond, } { *v = clampIntPtr(*v) } diff --git a/internal/agent/config.go b/internal/agent/config.go index ed43a7a..705a14e 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.3.1. +// contract v0.4.0. package agent import ( diff --git a/internal/agent/wire/wire.go b/internal/agent/wire/wire.go index b746cd2..0f8839d 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.2.1: the requests that the agent sends and the replies that it reads. +// v0.4.0: the requests that the agent sends and the replies that it reads. package wire import ( @@ -63,6 +63,17 @@ type Sample struct { MemoryPressureMaxPercent *float64 `json:"memory_pressure_max_percent"` IOPressurePercent *float64 `json:"io_pressure_percent"` IOPressureMaxPercent *float64 `json:"io_pressure_max_percent"` + + // The disk activity of the minute, and its peaks each second within the + // minute (contract v0.4.0). nil (JSON null) means "not known". + DiskReadBytes *uint64 `json:"disk_read_bytes"` + DiskWriteBytes *uint64 `json:"disk_write_bytes"` + DiskReadOps *uint64 `json:"disk_read_ops"` + DiskWriteOps *uint64 `json:"disk_write_ops"` + DiskReadMaxBytesPerSecond *uint64 `json:"disk_read_max_bytes_per_second"` + 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"` } // MetricsReply is the reply to POST /agent/v1/metrics. diff --git a/internal/metrics/disk.go b/internal/metrics/disk.go new file mode 100644 index 0000000..7f32787 --- /dev/null +++ b/internal/metrics/disk.go @@ -0,0 +1,98 @@ +package metrics + +import ( + "os" + "slices" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +// sectorSize is the unit of the sector counts in /proc/diskstats, for each +// disk, whatever the sector size of the disk. +const sectorSize = 512 + +// readDisks returns the counters of the disks that have a hardware device: a +// name in /sys/block with a device link, for example vda, sda or nvme0n1. +// Partitions are not in /sys/block, and loop, ram, zram, device mapper and +// software RAID disks have no device: their activity is already in a hardware +// disk. It returns nil when /proc/diskstats cannot be read. +func (c *Collector) readDisks() map[string]diskCounters { + all, err := parseFile(c, "proc/diskstats", parseDiskstats) + if err != nil { + c.log.Debug("reading the disk activity", "error", err) + return nil + } + + disks := map[string]diskCounters{} + for name, d := range all { + if _, err := os.Lstat(c.file("sys/block", name, "device")); err == nil { + disks[name] = d + } + } + return disks +} + +// diskDelta returns the activity between two readings. It adds the disks +// that both readings have, so a new or a removed disk makes no spike. ok is +// false when the activity is not known: a reading without disks, a reboot, a +// reading older than maxAge, or a counter that went back. +func diskDelta(a, b reading) (d diskCounters, ok bool) { + if len(a.Disks) == 0 || len(b.Disks) == 0 || a.BootID != b.BootID { + return diskCounters{}, false + } + if age := b.At.Sub(a.At); age <= 0 || age > maxAge { + return diskCounters{}, false + } + + for name, cur := range b.Disks { + prev, found := a.Disks[name] + if !found { + continue + } + if cur.ReadOps < prev.ReadOps || cur.ReadSectors < prev.ReadSectors || + cur.WriteOps < prev.WriteOps || cur.WriteSectors < prev.WriteSectors { + return diskCounters{}, false + } + d.ReadOps += cur.ReadOps - prev.ReadOps + d.ReadSectors += cur.ReadSectors - prev.ReadSectors + d.WriteOps += cur.WriteOps - prev.WriteOps + d.WriteSectors += cur.WriteSectors - prev.WriteSectors + } + + return d, true +} + +// setDiskActivity sets the disk activity of the minute from prev to cur in s, +// and the peaks of the windows in all. The eight fields stay nil when the +// activity of the minute is not known. The peaks also stay nil when a window +// has a counter that went back. +func setDiskActivity(s *wire.Sample, prev *reading, all []reading, cur reading) { + if prev == nil { + return + } + d, ok := diskDelta(*prev, cur) + if !ok { + return + } + readBytes, writeBytes := d.ReadSectors*sectorSize, d.WriteSectors*sectorSize + s.DiskReadBytes, s.DiskWriteBytes = &readBytes, &writeBytes + s.DiskReadOps, s.DiskWriteOps = &d.ReadOps, &d.WriteOps + + // The control plane reads the value of the minute as the value / 60. + peak := [4]uint64{readBytes / 60, writeBytes / 60, d.ReadOps / 60, d.WriteOps / 60} + // A reading whose disks were not read is left out: its windows join. + all = slices.DeleteFunc(slices.Clone(all), func(r reading) bool { return r.Disks == nil }) + for i := 1; i < len(all); i++ { + a, b := all[i-1], all[i] + w, ok := diskDelta(a, b) + if !ok { + return + } + secs := b.At.Sub(a.At).Seconds() + for j, v := range []uint64{w.ReadSectors * sectorSize, w.WriteSectors * sectorSize, w.ReadOps, w.WriteOps} { + peak[j] = max(peak[j], uint64(float64(v)/secs)) + } + } + s.DiskReadMaxBytesPerSecond, s.DiskWriteMaxBytesPerSecond = &peak[0], &peak[1] + s.DiskReadMaxOpsPerSecond, s.DiskWriteMaxOpsPerSecond = &peak[2], &peak[3] +} diff --git a/internal/metrics/disk_test.go b/internal/metrics/disk_test.go new file mode 100644 index 0000000..26951f3 --- /dev/null +++ b/internal/metrics/disk_test.go @@ -0,0 +1,309 @@ +package metrics + +import ( + "fmt" + "maps" + "os" + "path/filepath" + "slices" + "strings" + "testing" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +// A /proc/diskstats of a DigitalOcean server with Ubuntu 24.04: two disks, +// the partitions of vda, and the loop devices of snap. +const diskstats = ` 7 0 loop0 11 0 28 0 0 0 0 0 0 0 0 0 0 0 0 0 0 + 7 1 loop1 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 + 253 0 vda 33376 8239 2465787 22011 165725 77576 2598846 818732 0 202221 889331 2747 0 2299640 1263 22315 47324 + 253 1 vda1 32655 7544 2428229 21737 165698 77555 2598618 818558 0 211841 841559 2747 0 2299640 1263 0 0 + 253 14 vda14 217 0 1984 58 0 0 0 0 0 52 58 0 0 0 0 0 0 + 253 15 vda15 212 661 18510 80 2 0 2 3 0 52 83 0 0 0 0 0 0 + 259 0 vda16 192 34 13224 109 25 21 226 169 0 253 279 0 0 0 0 0 0 + 253 16 vdb 101 3 792 12 0 0 0 0 0 12 12 0 0 0 0 0 0 +` + +// blockDevice marks a disk in /sys/block as a hardware device. +func (s *server) blockDevice(name string) { + s.t.Helper() + if err := os.MkdirAll(filepath.Join(s.root, "sys/block", name, "device"), 0o755); err != nil { + s.t.Fatal(err) + } +} + +// disks writes /proc/diskstats with the counters of each disk, and a loop +// device and a partition with large counters. +func (s *server) disks(disks map[string]diskCounters) { + var b strings.Builder + fmt.Fprintf(&b, " 7 0 loop0 99999 0 99999 0 99999 0 99999 0 0 0 0 0 0 0 0 0 0\n") + for _, name := range slices.Sorted(maps.Keys(disks)) { + d := disks[name] + fmt.Fprintf(&b, " 253 0 %s %d 0 %d 0 %d 0 %d 0 0 0 0 0 0 0 0 0 0\n", name, d.ReadOps, d.ReadSectors, d.WriteOps, d.WriteSectors) + } + fmt.Fprintf(&b, " 253 1 vda1 99999 0 99999 0 99999 0 99999 0 0 0 0 0 0 0 0 0 0\n") + s.write("proc/diskstats", b.String()) +} + +func TestParseDiskstats(t *testing.T) { + all, err := parseDiskstats([]byte(diskstats)) + if err != nil { + t.Fatal(err) + } + want := diskCounters{ReadOps: 33376, ReadSectors: 2465787, WriteOps: 165725, WriteSectors: 2598846} + if all["vda"] != want { + t.Errorf("vda = %+v, want %+v", all["vda"], want) + } + if len(all) != 8 { + t.Errorf("%d lines, want 8", len(all)) + } + // A kernel before 4.18 has 14 columns: no discards and flushes. + all, err = parseDiskstats([]byte(" 8 0 sda 1 2 3 4 5 6 7 8 9 10 11\n")) + if err != nil || all["sda"] != (diskCounters{ReadOps: 1, ReadSectors: 3, WriteOps: 5, WriteSectors: 7}) { + t.Errorf("parseDiskstats(14 columns) = %+v, %v", all, err) + } + if _, err := parseDiskstats([]byte(" 8 0 sda 1 x 3 4 5 6 7 8 9 10 11\n")); err != nil { + t.Errorf("parseDiskstats() = %v, want no error: a column that is not used is not examined", err) + } + if _, err := parseDiskstats([]byte(" 8 0 sda x 2 3 4 5 6 7 8 9 10 11\n")); err == nil { + t.Error("parseDiskstats() = nil error, want an error for a bad count") + } +} + +func TestOnlyHardwareDisksAreCounted(t *testing.T) { + srv := newServer(t) + srv.write("proc/diskstats", diskstats) + srv.blockDevice("vda") + srv.blockDevice("vdb") + // A loop device is in /sys/block, but it has no device. + if err := os.MkdirAll(filepath.Join(srv.root, "sys/block/loop0"), 0o755); err != nil { + t.Fatal(err) + } + + got := slices.Sorted(maps.Keys(srv.collector(t.TempDir()).readDisks())) + if fmt.Sprint(got) != "[vda vdb]" { + t.Errorf("disks = %v, want [vda vdb]", got) + } +} + +// diskMinute takes a sample at base, five readings and the tick at base + +// 60 s, with the counters of vda at each of the seven readings. +func diskMinute(t *testing.T, srv *server, c *Collector, base time.Time, vda [7]diskCounters) wire.Sample { + t.Helper() + srv.disks(map[string]diskCounters{"vda": vda[0]}) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + for i := 1; i < 6; i++ { + srv.disks(map[string]diskCounters{"vda": vda[i]}) + c.Read(base.Add(time.Duration(i) * 10 * time.Second)) + } + srv.disks(map[string]diskCounters{"vda": vda[6]}) + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + return s +} + +func TestDiskActivity(t *testing.T) { + srv := newServer(t) + srv.blockDevice("vda") + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + // Each window reads 10 operations of 20 sectors, and writes 100 + // operations of 200 sectors. The window from 30 s to 40 s writes 5000 + // operations of 100000 sectors. + var vda [7]diskCounters + for i := 1; i < 7; i++ { + vda[i] = vda[i-1] + vda[i].ReadOps += 10 + vda[i].ReadSectors += 20 + if i == 4 { + vda[i].WriteOps += 5000 + vda[i].WriteSectors += 100000 + } else { + vda[i].WriteOps += 100 + vda[i].WriteSectors += 200 + } + } + s := diskMinute(t, srv, c, base, vda) + + for _, f := range []struct { + name string + got *uint64 + want uint64 + }{ + {"disk_read_bytes", s.DiskReadBytes, 120 * 512}, + {"disk_write_bytes", s.DiskWriteBytes, 101000 * 512}, + {"disk_read_ops", s.DiskReadOps, 60}, + {"disk_write_ops", s.DiskWriteOps, 5500}, + {"disk_read_max_bytes_per_second", s.DiskReadMaxBytesPerSecond, 20 * 512 / 10}, + {"disk_write_max_bytes_per_second", s.DiskWriteMaxBytesPerSecond, 100000 * 512 / 10}, + {"disk_read_max_ops_per_second", s.DiskReadMaxOpsPerSecond, 1}, + {"disk_write_max_ops_per_second", s.DiskWriteMaxOpsPerSecond, 500}, + } { + if f.got == nil || *f.got != f.want { + t.Errorf("%s = %v, want %d", f.name, ptr(f.got), f.want) + } + } + checkDiskOrder(t, s) +} + +func TestDiskActivityIsNotKnown(t *testing.T) { + tests := []struct { + name string + change func(*server) + after time.Duration + }{ + {"no hardware disk", func(s *server) { + if err := os.RemoveAll(filepath.Join(s.root, "sys/block")); err != nil { + s.t.Fatal(err) + } + }, time.Minute}, + {"no /proc/diskstats", func(s *server) { + if err := os.Remove(filepath.Join(s.root, "proc/diskstats")); err != nil { + s.t.Fatal(err) + } + }, time.Minute}, + {"reboot", func(s *server) { s.write("proc/sys/kernel/random/boot_id", "boot-2\n") }, time.Minute}, + {"counter went back", func(s *server) { s.disks(map[string]diskCounters{"vda": {ReadOps: 1}}) }, time.Minute}, + {"previous reading too old", func(*server) {}, 3 * time.Minute}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + srv := newServer(t) + srv.blockDevice("vda") + srv.disks(map[string]diskCounters{"vda": {ReadOps: 100, ReadSectors: 100, WriteOps: 100, WriteSectors: 100}}) + c := srv.collector(t.TempDir()) + now := time.Now().Add(time.Second) + if _, err := c.Sample(now); err != nil { + t.Fatal(err) + } + + tt.change(srv) + s, err := c.Sample(now.Add(tt.after)) + if err != nil { + t.Fatal(err) + } + checkNoDiskActivity(t, s) + }) + } +} + +func TestFirstSampleHasNoDiskActivity(t *testing.T) { + srv := newServer(t) + srv.blockDevice("vda") + srv.disks(map[string]diskCounters{"vda": {ReadOps: 100}}) + c := srv.collector(t.TempDir()) + srv.disks(map[string]diskCounters{"vda": {ReadOps: 200}}) + + s, err := c.Sample(time.Now().Add(30 * time.Second)) + if err != nil { + t.Fatal(err) + } + checkNoDiskActivity(t, s) +} + +func TestNewDiskMakesNoSpike(t *testing.T) { + srv := newServer(t) + srv.blockDevice("vda") + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + srv.disks(map[string]diskCounters{"vda": {ReadOps: 100}}) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + // A volume is attached at 30 s, with counters from its own start. + srv.blockDevice("sdb") + srv.disks(map[string]diskCounters{"vda": {ReadOps: 130}, "sdb": {ReadOps: 1 << 40}}) + c.Read(base.Add(30 * time.Second)) + srv.disks(map[string]diskCounters{"vda": {ReadOps: 160}, "sdb": {ReadOps: 1<<40 + 30}}) + + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.DiskReadOps == nil || *s.DiskReadOps != 60 { + t.Errorf("disk_read_ops = %v, want 60: the new sdb counts from the next sample", ptr(s.DiskReadOps)) + } + if s.DiskReadMaxOpsPerSecond == nil || *s.DiskReadMaxOpsPerSecond != 2 { + t.Errorf("disk_read_max_ops_per_second = %v, want 2: 30 of vda and 30 of sdb in 30 s", ptr(s.DiskReadMaxOpsPerSecond)) + } + checkDiskOrder(t, s) +} + +func TestAReadingWithoutDisksJoinsTheWindows(t *testing.T) { + srv := newServer(t) + srv.blockDevice("vda") + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + srv.disks(map[string]diskCounters{"vda": {}}) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + // /proc/diskstats cannot be read at 20 s. The other counters can. + if err := os.Remove(filepath.Join(srv.root, "proc/diskstats")); err != nil { + t.Fatal(err) + } + c.Read(base.Add(20 * time.Second)) + srv.disks(map[string]diskCounters{"vda": {WriteOps: 400}}) + c.Read(base.Add(40 * time.Second)) + srv.disks(map[string]diskCounters{"vda": {WriteOps: 600}}) + + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.DiskWriteMaxOpsPerSecond == nil || *s.DiskWriteMaxOpsPerSecond != 10 { + t.Errorf("disk_write_max_ops_per_second = %v, want 10: 400 in 40 s, or 200 in 20 s", ptr(s.DiskWriteMaxOpsPerSecond)) + } + + // Now the busy window is before the reading without disks: 600 in the + // joined window of 40 s is 15 each second. + if err := os.Remove(filepath.Join(srv.root, "proc/diskstats")); err != nil { + t.Fatal(err) + } + c.Read(base.Add(80 * time.Second)) + srv.disks(map[string]diskCounters{"vda": {WriteOps: 1200}}) + c.Read(base.Add(100 * time.Second)) + srv.disks(map[string]diskCounters{"vda": {WriteOps: 1200}}) + s, err = c.Sample(base.Add(2 * time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.DiskWriteMaxOpsPerSecond == nil || *s.DiskWriteMaxOpsPerSecond != 15 { + t.Errorf("disk_write_max_ops_per_second = %v, want 15 from the joined window of 60 s to 100 s", ptr(s.DiskWriteMaxOpsPerSecond)) + } +} + +func checkNoDiskActivity(t *testing.T, s wire.Sample) { + t.Helper() + for _, v := range []*uint64{s.DiskReadBytes, s.DiskWriteBytes, s.DiskReadOps, s.DiskWriteOps, + s.DiskReadMaxBytesPerSecond, s.DiskWriteMaxBytesPerSecond, s.DiskReadMaxOpsPerSecond, s.DiskWriteMaxOpsPerSecond} { + if v != nil { + t.Errorf("a disk activity field = %d, want null", *v) + } + } +} + +// checkDiskOrder checks that a value of the minute is at most 60 times its +// peak, up to rounding. +func checkDiskOrder(t *testing.T, s wire.Sample) { + t.Helper() + for _, p := range [][2]*uint64{ + {s.DiskReadBytes, s.DiskReadMaxBytesPerSecond}, + {s.DiskWriteBytes, s.DiskWriteMaxBytesPerSecond}, + {s.DiskReadOps, s.DiskReadMaxOpsPerSecond}, + {s.DiskWriteOps, s.DiskWriteMaxOpsPerSecond}, + } { + if p[0] == nil || p[1] == nil || *p[0] > 60**p[1]+59 { + t.Errorf("value %v, peak %v: want value ≤ 60 × peak", ptr(p[0]), ptr(p[1])) + } + } +} diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index b282c69..8b32869 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -1,7 +1,7 @@ // Package metrics measures a Linux server for the monitoring agent: CPU, -// load, memory, swap, disk and network each minute, with the peaks of the -// minute from a reading each 10 seconds, and the status of the server. It -// needs no root. +// 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. package metrics import ( @@ -53,6 +53,9 @@ type reading struct { Net map[string]netCounters `json:"net"` // PSI is nil when the kernel has no pressure information. PSI *psiTotals `json:"psi,omitempty"` + // Disks is nil when /proc/diskstats cannot be read, and empty when no + // disk has a hardware device. + Disks map[string]diskCounters `json:"disks,omitempty"` // mem is not saved: only the readings of the minute give its peak. mem memory } @@ -121,9 +124,10 @@ func New(root, stateDir string, log *slog.Logger) *Collector { c.prev = &saved c.readings = []reading{start} default: - // The traffic and the pressure since the start are not those of one - // minute: the first sample sends them as not known. - start.Net, start.PSI = nil, nil + // The traffic, the pressure and the disk activity since the start + // are not those of one minute: the first sample sends them as not + // known. + start.Net, start.PSI, start.Disks = nil, nil, nil c.prev = &start c.fromStart = true } @@ -277,7 +281,7 @@ func (c *Collector) read(now time.Time) (reading, error) { return reading{}, err } - return reading{BootID: string(bytes.TrimSpace(bootID)), At: now, CPU: cpu, Net: net, PSI: c.readPSI(), mem: mem}, nil + return reading{BootID: string(bytes.TrimSpace(bootID)), At: now, CPU: cpu, Net: net, PSI: c.readPSI(), Disks: c.readDisks(), mem: mem}, nil } // interfaces returns the network interfaces that have a hardware device and diff --git a/internal/metrics/parse.go b/internal/metrics/parse.go index 1e642e0..34512d4 100644 --- a/internal/metrics/parse.go +++ b/internal/metrics/parse.go @@ -228,3 +228,40 @@ func parsePSI(data []byte) (uint64, error) { return 0, fmt.Errorf("pressure: no total on a \"some\" line") } + +// diskCounters are the completed operations and the 512-byte sectors of one +// disk, from /proc/diskstats. +type diskCounters struct { + ReadOps uint64 `json:"read_ops"` + ReadSectors uint64 `json:"read_sectors"` + WriteOps uint64 `json:"write_ops"` + WriteSectors uint64 `json:"write_sectors"` +} + +// parseDiskstats reads /proc/diskstats. Each line is +// +// major minor name reads merged sectors ms writes merged sectors ms ... +// +// The reads and writes completed are columns 4 and 8, and the sectors read +// and written are columns 6 and 10. A sector is 512 bytes for each disk. +func parseDiskstats(data []byte) (map[string]diskCounters, error) { + out := map[string]diskCounters{} + s := bufio.NewScanner(bytes.NewReader(data)) + for s.Scan() { + fields := strings.Fields(s.Text()) + if len(fields) < 10 { + continue + } + var v [4]uint64 + for i, col := range []int{3, 5, 7, 9} { + n, err := strconv.ParseUint(fields[col], 10, 64) + if err != nil { + return nil, fmt.Errorf("/proc/diskstats: bad counters for %s", fields[2]) + } + v[i] = n + } + out[fields[2]] = diskCounters{ReadOps: v[0], ReadSectors: v[1], WriteOps: v[2], WriteSectors: v[3]} + } + + return out, nil +} diff --git a/internal/metrics/windows.go b/internal/metrics/windows.go index a11dbd5..974f78c 100644 --- a/internal/metrics/windows.go +++ b/internal/metrics/windows.go @@ -33,10 +33,10 @@ func chain(prev *reading, readings []reading, cur reading) []reading { return append(out, cur) } -// setPeaks sets the peaks of the minute in s, and the pressure: the highest -// value of the windows from one reading to the next (contract v0.4.0). s already holds -// the values of the minute, from prev to cur. A peak is never less than the -// value of its minute. +// setPeaks sets the peaks of the minute in s, the pressure and the disk +// activity. A peak is the highest value of the windows from one reading to +// the next (contract v0.4.0). s already holds the values of the minute, from +// prev to cur. A peak is never less than the value of its minute. func setPeaks(s *wire.Sample, prev *reading, readings []reading, cur reading) { all := chain(prev, readings, cur) @@ -61,6 +61,7 @@ func setPeaks(s *wire.Sample, prev *reading, readings []reading, cur reading) { } setPressure(s, prev, all, cur) + setDiskActivity(s, prev, all, cur) if !s.NetCountersReset { if in, out, ok := netPeaks(all); ok { From 1755fcf73f9a8f1a8cab26a71a25f8936975704f Mon Sep 17 00:00:00 2001 From: nabil1440 <52530910+nabil1440@users.noreply.github.com> Date: Thu, 24 Sep 2026 12:01:09 +0600 Subject: [PATCH 2/2] fix(agent): send no disk activity when a counter went back within the minute The contract makes all eight disk fields null when a counter went back. A disk that is attached again under the same name can pass the check of the minute and still give a wrong delta; now a window that went back nulls the values of the minute too. A reading with no disks for a moment joins its windows. --- internal/metrics/disk.go | 13 +++++++------ internal/metrics/disk_test.go | 25 +++++++++++++++++++++++++ 2 files changed, 32 insertions(+), 6 deletions(-) diff --git a/internal/metrics/disk.go b/internal/metrics/disk.go index 7f32787..f5c88f8 100644 --- a/internal/metrics/disk.go +++ b/internal/metrics/disk.go @@ -64,8 +64,8 @@ func diskDelta(a, b reading) (d diskCounters, ok bool) { // setDiskActivity sets the disk activity of the minute from prev to cur in s, // and the peaks of the windows in all. The eight fields stay nil when the -// activity of the minute is not known. The peaks also stay nil when a window -// has a counter that went back. +// activity of the minute is not known, or when a window has a counter that +// went back: then the delta of the minute is not the real activity either. func setDiskActivity(s *wire.Sample, prev *reading, all []reading, cur reading) { if prev == nil { return @@ -75,13 +75,11 @@ func setDiskActivity(s *wire.Sample, prev *reading, all []reading, cur reading) return } readBytes, writeBytes := d.ReadSectors*sectorSize, d.WriteSectors*sectorSize - s.DiskReadBytes, s.DiskWriteBytes = &readBytes, &writeBytes - s.DiskReadOps, s.DiskWriteOps = &d.ReadOps, &d.WriteOps // The control plane reads the value of the minute as the value / 60. peak := [4]uint64{readBytes / 60, writeBytes / 60, d.ReadOps / 60, d.WriteOps / 60} - // A reading whose disks were not read is left out: its windows join. - all = slices.DeleteFunc(slices.Clone(all), func(r reading) bool { return r.Disks == nil }) + // A reading without disks is left out: its windows join. + all = slices.DeleteFunc(slices.Clone(all), func(r reading) bool { return len(r.Disks) == 0 }) for i := 1; i < len(all); i++ { a, b := all[i-1], all[i] w, ok := diskDelta(a, b) @@ -93,6 +91,9 @@ func setDiskActivity(s *wire.Sample, prev *reading, all []reading, cur reading) peak[j] = max(peak[j], uint64(float64(v)/secs)) } } + + s.DiskReadBytes, s.DiskWriteBytes = &readBytes, &writeBytes + s.DiskReadOps, s.DiskWriteOps = &d.ReadOps, &d.WriteOps s.DiskReadMaxBytesPerSecond, s.DiskWriteMaxBytesPerSecond = &peak[0], &peak[1] s.DiskReadMaxOpsPerSecond, s.DiskWriteMaxOpsPerSecond = &peak[2], &peak[3] } diff --git a/internal/metrics/disk_test.go b/internal/metrics/disk_test.go index 26951f3..e155d12 100644 --- a/internal/metrics/disk_test.go +++ b/internal/metrics/disk_test.go @@ -194,6 +194,31 @@ func TestDiskActivityIsNotKnown(t *testing.T) { } } +func TestACounterThatWentBackInAWindowMakesTheMinuteNotKnown(t *testing.T) { + srv := newServer(t) + srv.blockDevice("vda") + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + // vda is attached again at 20 s, with the same name: its counter + // starts from 0. At the tick it is above the tick before again, but the + // delta of the minute is not the real activity. + srv.disks(map[string]diskCounters{"vda": {ReadOps: 1000}}) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + c.Read(base.Add(10 * time.Second)) + srv.disks(map[string]diskCounters{"vda": {ReadOps: 5}}) + c.Read(base.Add(20 * time.Second)) + srv.disks(map[string]diskCounters{"vda": {ReadOps: 2000}}) + + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + checkNoDiskActivity(t, s) +} + func TestFirstSampleHasNoDiskActivity(t *testing.T) { srv := newServer(t) srv.blockDevice("vda")