Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion internal/agent/clean.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
} {
Expand Down
10 changes: 10 additions & 0 deletions internal/agent/wire/wire.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
13 changes: 9 additions & 4 deletions internal/metrics/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand All @@ -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
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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
Expand Down
26 changes: 26 additions & 0 deletions internal/metrics/parse.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
86 changes: 86 additions & 0 deletions internal/metrics/pressure.go
Original file line number Diff line number Diff line change
@@ -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
}
197 changes: 197 additions & 0 deletions internal/metrics/pressure_test.go
Original file line number Diff line number Diff line change
@@ -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))
}
}
Loading
Loading