diff --git a/internal/agent/clean.go b/internal/agent/clean.go index d7b394e..a1dceae 100644 --- a/internal/agent/clean.go +++ b/internal/agent/clean.go @@ -32,7 +32,12 @@ func cleanSample(s wire.Sample) wire.Sample { } { *v = clampInt(*v) } - s.CPUMaxPercent = clampPtr(s.CPUMaxPercent, 0, 100) + for _, v := range []**float64{ + &s.CPUMaxPercent, &s.CPUPressurePercent, &s.CPUPressureMaxPercent, &s.MemoryPressurePercent, + &s.MemoryPressureMaxPercent, &s.IOPressurePercent, &s.IOPressureMaxPercent, + } { + *v = clampPtr(*v, 0, 100) + } for _, v := range []**uint64{ &s.MemoryUsedMaxBytes, &s.SwapUsedMaxBytes, &s.NetInMaxBytesPerSecond, &s.NetOutMaxBytesPerSecond, } { diff --git a/internal/agent/wire/wire.go b/internal/agent/wire/wire.go index 32e4673..b746cd2 100644 --- a/internal/agent/wire/wire.go +++ b/internal/agent/wire/wire.go @@ -53,6 +53,16 @@ type Sample struct { SwapUsedMaxBytes *uint64 `json:"swap_used_max_bytes"` NetInMaxBytesPerSecond *uint64 `json:"net_in_max_bytes_per_second"` NetOutMaxBytesPerSecond *uint64 `json:"net_out_max_bytes_per_second"` + + // The share of the time in which at least one task waited for the CPU, + // the memory or the disk (PSI), 0 to 100, and its peak within the minute + // (contract v0.4.0). nil (JSON null) means "not known". + CPUPressurePercent *float64 `json:"cpu_pressure_percent"` + CPUPressureMaxPercent *float64 `json:"cpu_pressure_max_percent"` + MemoryPressurePercent *float64 `json:"memory_pressure_percent"` + MemoryPressureMaxPercent *float64 `json:"memory_pressure_max_percent"` + IOPressurePercent *float64 `json:"io_pressure_percent"` + IOPressureMaxPercent *float64 `json:"io_pressure_max_percent"` } // MetricsReply is the reply to POST /agent/v1/metrics. diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index 1638365..b282c69 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -51,6 +51,8 @@ type reading struct { At time.Time `json:"at"` CPU cpuTimes `json:"cpu"` Net map[string]netCounters `json:"net"` + // PSI is nil when the kernel has no pressure information. + PSI *psiTotals `json:"psi,omitempty"` // mem is not saved: only the readings of the minute give its peak. mem memory } @@ -76,8 +78,10 @@ type Collector struct { // tests set it to 0. fromStart bool minFirst time.Duration - // noInterface is true after the warning that no interface is counted. + // noInterface is true after the warning that no interface is counted, + // and noPSI after the warning that the kernel has no PSI. noInterface bool + noPSI bool updatesAt time.Time updatesKnown bool @@ -117,8 +121,9 @@ func New(root, stateDir string, log *slog.Logger) *Collector { c.prev = &saved c.readings = []reading{start} default: - // The traffic since the start is not the traffic of one minute. - start.Net = nil + // 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 c.prev = &start c.fromStart = true } @@ -272,7 +277,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, mem: mem}, nil + return reading{BootID: string(bytes.TrimSpace(bootID)), At: now, CPU: cpu, Net: net, PSI: c.readPSI(), 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 d5f6724..1e642e0 100644 --- a/internal/metrics/parse.go +++ b/internal/metrics/parse.go @@ -202,3 +202,29 @@ func parseAptCheck(out []byte) (total, security uint64, err error) { return total, security, nil } + +// parsePSI reads the total of the "some" line of a /proc/pressure file: the +// microseconds in which at least one task waited for the resource. +// +// some avg10=0.00 avg60=0.00 avg300=0.00 total=12345 +// full avg10=0.00 avg60=0.00 avg300=0.00 total=0 +func parsePSI(data []byte) (uint64, error) { + s := bufio.NewScanner(bytes.NewReader(data)) + for s.Scan() { + fields := strings.Fields(s.Text()) + if len(fields) == 0 || fields[0] != "some" { + continue + } + for _, f := range fields[1:] { + if v, ok := strings.CutPrefix(f, "total="); ok { + n, err := strconv.ParseUint(v, 10, 64) + if err != nil { + return 0, fmt.Errorf("pressure: bad total %q", v) + } + return n, nil + } + } + } + + return 0, fmt.Errorf("pressure: no total on a \"some\" line") +} diff --git a/internal/metrics/pressure.go b/internal/metrics/pressure.go new file mode 100644 index 0000000..e07f15b --- /dev/null +++ b/internal/metrics/pressure.go @@ -0,0 +1,86 @@ +package metrics + +import ( + "slices" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +// psiTotals are the "some" totals of /proc/pressure, in microseconds. +type psiTotals struct { + CPU uint64 `json:"cpu"` + Memory uint64 `json:"memory"` + IO uint64 `json:"io"` +} + +// readPSI reads the three pressure files. It returns nil when the kernel has +// no PSI: no /proc/pressure, or a kernel booted with psi=0, where a read +// fails. It logs this one time. +func (c *Collector) readPSI() *psiTotals { + var t psiTotals + for _, f := range []struct { + name string + v *uint64 + }{{"cpu", &t.CPU}, {"memory", &t.Memory}, {"io", &t.IO}} { + n, err := parseFile(c, "proc/pressure/"+f.name, parsePSI) + if err != nil { + if !c.noPSI { + c.noPSI = true + c.log.Warn("no pressure (PSI) to send: the kernel has no PSI, or it is off", "error", err) + } + return nil + } + *f.v = n + } + + return &t +} + +// pressure returns the share of the time between two readings in which at +// least one task waited for the CPU, the memory and the disk, 0 to 100. ok is +// false when the share is not known: a reading without PSI, a reboot, a +// reading older than maxAge, or a counter that went back. +func pressure(a, b reading) (cpu, mem, io float64, ok bool) { + if a.PSI == nil || b.PSI == nil || a.BootID != b.BootID { + return 0, 0, 0, false + } + d := b.At.Sub(a.At) + if d <= 0 || d > maxAge { + return 0, 0, 0, false + } + if b.PSI.CPU < a.PSI.CPU || b.PSI.Memory < a.PSI.Memory || b.PSI.IO < a.PSI.IO { + return 0, 0, 0, false + } + + us := float64(d.Microseconds()) + share := func(from, to uint64) float64 { + return min(float64(to-from)/us*100, 100) + } + return share(a.PSI.CPU, b.PSI.CPU), share(a.PSI.Memory, b.PSI.Memory), share(a.PSI.IO, b.PSI.IO), true +} + +// setPressure sets the pressure of the minute from prev to cur in s, and the +// peak of the windows in all. The six fields stay nil when the pressure of the +// minute is not known. +func setPressure(s *wire.Sample, prev *reading, all []reading, cur reading) { + if prev == nil { + return + } + cpu, mem, io, ok := pressure(*prev, cur) + if !ok { + return + } + + cpuMax, memMax, ioMax := cpu, mem, io + // A reading without PSI is left out: its windows join. + all = slices.DeleteFunc(slices.Clone(all), func(r reading) bool { return r.PSI == nil }) + for i := 1; i < len(all); i++ { + if c, m, o, ok := pressure(all[i-1], all[i]); ok { + cpuMax, memMax, ioMax = max(cpuMax, c), max(memMax, m), max(ioMax, o) + } + } + + s.CPUPressurePercent, s.CPUPressureMaxPercent = &cpu, &cpuMax + s.MemoryPressurePercent, s.MemoryPressureMaxPercent = &mem, &memMax + s.IOPressurePercent, s.IOPressureMaxPercent = &io, &ioMax +} diff --git a/internal/metrics/pressure_test.go b/internal/metrics/pressure_test.go new file mode 100644 index 0000000..3e698e5 --- /dev/null +++ b/internal/metrics/pressure_test.go @@ -0,0 +1,197 @@ +package metrics + +import ( + "fmt" + "os" + "path/filepath" + "testing" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +// psi writes the three /proc/pressure files with the "some" totals, in +// microseconds. +func (s *server) psi(cpu, memory, io uint64) { + for name, total := range map[string]uint64{"cpu": cpu, "memory": memory, "io": io} { + s.write("proc/pressure/"+name, fmt.Sprintf("some avg10=1.00 avg60=2.00 avg300=3.00 total=%d\nfull avg10=0.00 avg60=0.00 avg300=0.00 total=7\n", total)) + } +} + +func TestParsePSI(t *testing.T) { + n, err := parsePSI([]byte("some avg10=0.12 avg60=0.34 avg300=0.56 total=987654321\nfull avg10=0.00 avg60=0.00 avg300=0.00 total=5\n")) + if err != nil || n != 987654321 { + t.Errorf("parsePSI() = %d, %v; want the total of the some line", n, err) + } + // The CPU file of an older kernel has no full line. + if n, err := parsePSI([]byte("some avg10=0.00 avg60=0.00 avg300=0.00 total=42\n")); err != nil || n != 42 { + t.Errorf("parsePSI() = %d, %v; want 42", n, err) + } + for _, bad := range []string{"", "full avg10=0.00 total=5\n", "some avg10=0.00\n", "some total=x\n"} { + if _, err := parsePSI([]byte(bad)); err == nil { + t.Errorf("parsePSI(%q) = nil error, want an error", bad) + } + } +} + +// pressureMinute takes a sample at base, five readings and the tick at +// base + 60 s. totals are the CPU, memory and I/O totals at each of the seven +// readings. +func pressureMinute(t *testing.T, srv *server, c *Collector, base time.Time, totals [7][3]uint64) wire.Sample { + t.Helper() + srv.psi(totals[0][0], totals[0][1], totals[0][2]) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + for i := 1; i < 6; i++ { + srv.psi(totals[i][0], totals[i][1], totals[i][2]) + c.Read(base.Add(time.Duration(i) * 10 * time.Second)) + } + srv.psi(totals[6][0], totals[6][1], totals[6][2]) + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + return s +} + +func TestPressure(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + // The CPU waits 5 s in the window from 20 s to 30 s: 50% of that window, + // and 6 s of the minute: 10%. The memory never waits. The I/O waits + // 0.6 s in each window: 1%. + s := pressureMinute(t, srv, c, base, [7][3]uint64{ + {1000, 0, 0}, + {201000, 0, 100000}, + {401000, 0, 200000}, + {5401000, 0, 300000}, + {5601000, 0, 400000}, + {5801000, 0, 500000}, + {6001000, 0, 600000}, + }) + + want := map[string][2]float64{"cpu": {10, 50}, "memory": {0, 0}, "io": {1, 1}} + got := map[string][2]*float64{ + "cpu": {s.CPUPressurePercent, s.CPUPressureMaxPercent}, + "memory": {s.MemoryPressurePercent, s.MemoryPressureMaxPercent}, + "io": {s.IOPressurePercent, s.IOPressureMaxPercent}, + } + for name, w := range want { + g := got[name] + if g[0] == nil || g[1] == nil || !near(*g[0], w[0]) || !near(*g[1], w[1]) { + t.Errorf("%s pressure = %v, max %v; want %v, max %v", name, ptr(g[0]), ptr(g[1]), w[0], w[1]) + } + } +} + +func TestPressureIsNotKnown(t *testing.T) { + tests := []struct { + name string + change func(*server) + after time.Duration + }{ + {"no /proc/pressure", func(s *server) { + if err := os.RemoveAll(filepath.Join(s.root, "proc/pressure")); 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.psi(10, 10, 10) }, 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.psi(1000, 1000, 1000) + 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) + } + checkNoPressure(t, s) + }) + } +} + +func TestFirstSampleHasNoPressure(t *testing.T) { + srv := newServer(t) + srv.psi(1000, 1000, 1000) + c := srv.collector(t.TempDir()) + srv.psi(2000, 2000, 2000) + + s, err := c.Sample(time.Now().Add(30 * time.Second)) + if err != nil { + t.Fatal(err) + } + checkNoPressure(t, s) +} + +func TestPressureContinuesAfterARestart(t *testing.T) { + srv := newServer(t) + srv.psi(1000, 1000, 1000) + state := t.TempDir() + now := time.Now() + if _, err := srv.collector(state).Sample(now.Add(-30 * time.Second)); err != nil { + t.Fatal(err) + } + + c := srv.collector(state) + srv.psi(601000, 1000, 1000) + s, err := c.Sample(now.Add(30 * time.Second)) + if err != nil { + t.Fatal(err) + } + if s.CPUPressurePercent == nil || !near(*s.CPUPressurePercent, 1) { + t.Errorf("cpu_pressure_percent = %v, want 1 from the saved reading", ptr(s.CPUPressurePercent)) + } +} + +func checkNoPressure(t *testing.T, s wire.Sample) { + t.Helper() + for _, v := range []*float64{s.CPUPressurePercent, s.CPUPressureMaxPercent, s.MemoryPressurePercent, s.MemoryPressureMaxPercent, s.IOPressurePercent, s.IOPressureMaxPercent} { + if v != nil { + t.Errorf("a pressure field = %v, want null", *v) + } + } +} + +// near reports whether a is b, up to the rounding of the microseconds. +func near(a, b float64) bool { + return a > b-0.001 && a < b+0.001 +} + +func TestAReadingWithoutPressureJoinsTheWindows(t *testing.T) { + srv := newServer(t) + srv.psi(0, 0, 0) + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + + // The CPU waits 8 s from 0 s to 20 s, but the reading at 10 s has no + // PSI: the peak is 8 s of the joined window of 20 s. + srv.write("proc/pressure/cpu", "") + c.Read(base.Add(10 * time.Second)) + srv.psi(8000000, 0, 0) + c.Read(base.Add(20 * time.Second)) + srv.psi(9200000, 0, 0) + + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.CPUPressureMaxPercent == nil || !near(*s.CPUPressureMaxPercent, 40) { + t.Errorf("cpu_pressure_max_percent = %v, want 40 from the joined window of 0 s to 20 s", ptr(s.CPUPressureMaxPercent)) + } +} diff --git a/internal/metrics/windows.go b/internal/metrics/windows.go index 5938fb9..a11dbd5 100644 --- a/internal/metrics/windows.go +++ b/internal/metrics/windows.go @@ -33,8 +33,8 @@ func chain(prev *reading, readings []reading, cur reading) []reading { return append(out, cur) } -// setPeaks sets the peaks of the minute in s: the highest value of the -// windows from one reading to the next (contract v0.4.0). s already holds +// 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. func setPeaks(s *wire.Sample, prev *reading, readings []reading, cur reading) { @@ -60,6 +60,8 @@ func setPeaks(s *wire.Sample, prev *reading, readings []reading, cur reading) { s.CPUMaxPercent = &cpu } + setPressure(s, prev, all, cur) + if !s.NetCountersReset { if in, out, ok := netPeaks(all); ok { // The control plane reads the value of the minute as bytes / 60.