diff --git a/internal/agent/agent.go b/internal/agent/agent.go index b5a23e0..ffce100 100644 --- a/internal/agent/agent.go +++ b/internal/agent/agent.go @@ -31,9 +31,16 @@ type ControlPlane interface { PollCommands(ctx context.Context) (*wire.CommandsReply, error) } +// ErrNoSample is the error of a minute that has no sample, for example a +// first minute that is too short to measure. It is not a problem. +var ErrNoSample = errors.New("no sample for this minute") + // Collector measures the server. type Collector interface { - // Sample measures the minute that ends at now. + // Read takes a reading between two ticks, for the peaks of the minute. + Read(now time.Time) + // Sample measures the minute that ends at now. It returns an error that + // wraps ErrNoSample when the minute has no sample. Sample(now time.Time) (wire.Sample, error) // Status describes the server now. Status(ctx context.Context) wire.Status @@ -133,17 +140,36 @@ func run(ctx context.Context, cfg Config, log *slog.Logger, cp ControlPlane, col return nil } -// loop calls tick at the offset second of each minute until ctx is done. It -// returns true when a command ends the process. +// loop calls tick at the offset second of each minute until ctx is done. +// Between two ticks, it reads the server each 10 seconds for the peaks of the +// minute. It returns true when a command ends the process. func (a *agent) loop(ctx context.Context) (exit bool) { for { - next := nextAfter(time.Now(), a.last, a.cfg.Offset()) - timer := time.NewTimer(time.Until(next)) + now := time.Now() + next := nextAfter(now, a.last, a.cfg.Offset()) + wake := next + if a.collector != nil { + if r := nextReading(now, a.cfg.Offset()); r.Before(next) { + wake = r + } + } + + timer := time.NewTimer(time.Until(wake)) select { case <-ctx.Done(): timer.Stop() return false case <-timer.C: + if wake.Before(next) { + // The reading has the time at which it ran, so that a late + // timer or a step of the wall clock does not change the + // length of a window. + a.collector.Read(time.Now()) + // A reading that ran late must not skip the tick after it. + if time.Now().Before(next) { + continue + } + } a.last = next if a.tick(ctx, next) { return true @@ -152,6 +178,22 @@ func (a *agent) loop(ctx context.Context) (exit bool) { } } +// readEvery is the time between two readings of the server (contract v0.4.0). +const readEvery = 10 * time.Second + +// nextReading returns the next reading after now: offset past a full minute, +// and each 10 seconds after it. A tick is also a reading time: the loop then +// runs the tick instead. After a step back of the wall clock, a reading can +// come again; the collector ignores a reading that is not newer than its last. +func nextReading(now time.Time, offset time.Duration) time.Time { + t := now.Truncate(readEvery).Add(offset % readEvery) + for !t.After(now) { + t = t.Add(readEvery) + } + + return t +} + // maxStepBack is the largest step back of the wall clock after which the loop // still skips the minute that already ran. A larger step means that the clock // was wrong before (for example a VM that booted with its clock ahead, then @@ -182,10 +224,15 @@ func (a *agent) tick(ctx context.Context, now time.Time) (exit bool) { a.log.Debug("tick", "at", now) if a.collector != nil { - s, err := a.collector.Sample(now) - if err != nil { + // The reading has the time at which it ran; the sample has the time + // of its tick, for its minute on the control plane. + s, err := a.collector.Sample(time.Now()) + switch { + case errors.Is(err, ErrNoSample): + a.log.Info("no sample for this minute", "reason", err) + case err != nil: a.log.Warn("skipping the sample of this minute", "error", err) - } else { + default: s.RecordedAt = now a.outbox.addSample(cleanSample(s)) } diff --git a/internal/agent/agent_test.go b/internal/agent/agent_test.go index 61a6d4c..0f3a81a 100644 --- a/internal/agent/agent_test.go +++ b/internal/agent/agent_test.go @@ -186,6 +186,97 @@ func TestLoopTicksAtTheOffsetAndReportsEachInterval(t *testing.T) { }) } +func TestNextReading(t *testing.T) { + base := time.Date(2026, 9, 22, 10, 0, 0, 0, time.UTC) + tests := []struct { + now time.Time + offset time.Duration + want time.Time + }{ + {base, 17 * time.Second, base.Add(7 * time.Second)}, + {base.Add(7 * time.Second), 17 * time.Second, base.Add(17 * time.Second)}, + {base.Add(18 * time.Second), 17 * time.Second, base.Add(27 * time.Second)}, + {base.Add(58 * time.Second), 17 * time.Second, base.Add(67 * time.Second)}, + {base.Add(59*time.Second + 999*time.Millisecond), 0, base.Add(time.Minute)}, + {base.Add(3 * time.Second), 43 * time.Second, base.Add(3 * time.Second).Add(10 * time.Second)}, + } + + for _, tt := range tests { + if got := nextReading(tt.now, tt.offset); !got.Equal(tt.want) { + t.Errorf("nextReading(%s, %v) = %s, want %s", tt.now.Format(time.TimeOnly), tt.offset, got.Format(time.TimeOnly), tt.want.Format(time.TimeOnly)) + } + } +} + +func TestLoopReadsEach10SecondsBetweenTheTicks(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + collector := &fakeCollector{} + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error) + go func() { + done <- run(ctx, Config{ServerID: 17, StateDir: t.TempDir()}, slog.New(slog.DiscardHandler), &fakeCP{}, collector) + }() + + time.Sleep(2 * time.Minute) + cancel() + if err := <-done; err != nil { + t.Fatal(err) + } + + // The ticks are at :17. The readings are at :07, :27, :37, :47 and :57: + // the loop never reads at a tick, because the tick takes the reading. + start := time.Date(2000, 1, 1, 0, 0, 0, 0, time.UTC) + var want []time.Time + for s := 7 * time.Second; s < 2*time.Minute; s += 10 * time.Second { + if s%time.Minute != 17*time.Second { + want = append(want, start.Add(s)) + } + } + got := collector.readTimes() + if len(got) != len(want) { + t.Fatalf("readings at %v, want %v", got, want) + } + for i := range want { + if !got[i].Equal(want[i]) { + t.Errorf("reading %d at %s, want %s", i, got[i].Format(time.TimeOnly), want[i].Format(time.TimeOnly)) + } + } + if collector.n != 2 { + t.Errorf("%d samples, want 2 (at 0:17 and 1:17)", collector.n) + } + }) +} + +func TestASlowReadingDoesNotSkipTheTick(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + // Each reading takes 10 s: the reading at :07 ends at the tick. + collector := &fakeCollector{readTime: 10 * time.Second} + rec := &recorder{} + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error) + go func() { + done <- run(ctx, Config{ServerID: 17, StateDir: t.TempDir()}, slog.New(rec), &fakeCP{}, collector) + }() + + time.Sleep(2*time.Minute + 30*time.Second) + cancel() + if err := <-done; err != nil { + t.Fatal(err) + } + + start := time.Date(2000, 1, 1, 0, 0, 0, 0, time.UTC) + ticks := rec.times("tick") + if len(ticks) != 3 { + t.Fatalf("ticks at %v, want 3 ticks", ticks) + } + for i, got := range ticks { + if want := start.Add(time.Duration(i)*time.Minute + 17*time.Second); !got.Equal(want) { + t.Errorf("tick %d at %s, want %s", i, got.Format(time.TimeOnly), want.Format(time.TimeOnly)) + } + } + }) +} + func TestRunRefusesASecondAgent(t *testing.T) { dir := t.TempDir() unlock, err := lock(dir) diff --git a/internal/agent/clean.go b/internal/agent/clean.go index 6a9702b..d7b394e 100644 --- a/internal/agent/clean.go +++ b/internal/agent/clean.go @@ -32,6 +32,12 @@ func cleanSample(s wire.Sample) wire.Sample { } { *v = clampInt(*v) } + s.CPUMaxPercent = clampPtr(s.CPUMaxPercent, 0, 100) + for _, v := range []**uint64{ + &s.MemoryUsedMaxBytes, &s.SwapUsedMaxBytes, &s.NetInMaxBytesPerSecond, &s.NetOutMaxBytesPerSecond, + } { + *v = clampIntPtr(*v) + } return s } @@ -87,6 +93,15 @@ func clamp(v, lo, hi float64) float64 { return min(max(v, lo), hi) } +// clampPtr is clamp for a value that can be "not known" (nil). +func clampPtr(v *float64, lo, hi float64) *float64 { + if v == nil { + return nil + } + c := clamp(*v, lo, hi) + return &c +} + // truncate cuts s to at most n characters. The control plane counts // characters, not bytes. func truncate(s string, n int) string { diff --git a/internal/agent/fakes_test.go b/internal/agent/fakes_test.go index 2e11734..1f647f7 100644 --- a/internal/agent/fakes_test.go +++ b/internal/agent/fakes_test.go @@ -122,8 +122,26 @@ func (f *fakeCP) sampleCounts() []int { type fakeCollector struct { mu sync.Mutex n int + reads []time.Time status wire.Status err error + // readTime is the time that each reading takes. + readTime time.Duration +} + +func (c *fakeCollector) Read(now time.Time) { + c.mu.Lock() + c.reads = append(c.reads, now) + d := c.readTime + c.mu.Unlock() + time.Sleep(d) +} + +// readTimes returns the times of the readings between the ticks. +func (c *fakeCollector) readTimes() []time.Time { + c.mu.Lock() + defer c.mu.Unlock() + return slices.Clone(c.reads) } func (c *fakeCollector) Sample(time.Time) (wire.Sample, error) { diff --git a/internal/agent/outbox_test.go b/internal/agent/outbox_test.go index 0446e92..15da7e9 100644 --- a/internal/agent/outbox_test.go +++ b/internal/agent/outbox_test.go @@ -174,6 +174,45 @@ func TestCleanKeepsIntegersInTheRangeOfPHP(t *testing.T) { } } +func TestCleanPeaks(t *testing.T) { + cpu, big := 120.0, uint64(math.MaxUint64) + s := cleanSample(wire.Sample{CPUMaxPercent: &cpu, MemoryUsedMaxBytes: &big, NetInMaxBytesPerSecond: &big}) + if *s.CPUMaxPercent != 100 || *s.MemoryUsedMaxBytes != math.MaxInt64 || *s.NetInMaxBytesPerSecond != math.MaxInt64 { + t.Errorf("cleanSample() = %v, %d, %d; want 100 and the largest PHP integer", *s.CPUMaxPercent, *s.MemoryUsedMaxBytes, *s.NetInMaxBytesPerSecond) + } + if cpu != 120 { + t.Error("cleanSample() changed the value of the caller") + } + if s.SwapUsedMaxBytes != nil || s.NetOutMaxBytesPerSecond != nil { + t.Error("cleanSample() gave a value to a peak that is not known") + } +} + +// A sample that an older agent queued has no peaks. After an update, the new +// agent sends them as null: not known. +func TestQueuedSampleOfAnOlderAgentSendsNullPeaks(t *testing.T) { + dir := t.TempDir() + old := `[{"recorded_at":"2026-09-24T10:00:17Z","cpu_percent":12.5,"load_1":0.4,"memory_used_bytes":1,"memory_total_bytes":2,` + + `"swap_used_bytes":0,"swap_total_bytes":0,"disk_used_bytes":1,"disk_total_bytes":2,"net_in_bytes":5,"net_out_bytes":6,"net_counters_reset":false}]` + if err := os.WriteFile(filepath.Join(dir, "samples.json"), []byte(old), 0o600); err != nil { + t.Fatal(err) + } + + o := loadOutbox(dir, slog.New(slog.DiscardHandler)) + if len(o.samples) != 1 { + t.Fatalf("samples = %d, want the queued sample", len(o.samples)) + } + data, err := json.Marshal(o.samples[0]) + if err != nil { + t.Fatal(err) + } + for _, field := range []string{"cpu_max_percent", "memory_used_max_bytes", "swap_used_max_bytes", "net_in_max_bytes_per_second", "net_out_max_bytes_per_second"} { + if !strings.Contains(string(data), `"`+field+`":null`) { + t.Errorf("sample JSON = %s, want %s as null", data, field) + } + } +} + func TestCleanEventDropsACommandIDThatIsNotAULID(t *testing.T) { if e := cleanEvent(wire.Event{CommandID: "not-a-ulid"}); e.CommandID != "" { t.Errorf("command_id = %q, want it removed", e.CommandID) diff --git a/internal/agent/send_test.go b/internal/agent/send_test.go index d7a1540..5c3e0ed 100644 --- a/internal/agent/send_test.go +++ b/internal/agent/send_test.go @@ -319,6 +319,25 @@ func TestARetryDoesNotWaitForTheNextInterval(t *testing.T) { }) } +func TestAMinuteWithoutASampleIsNotAWarning(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + collector := &fakeCollector{err: fmt.Errorf("%w: the first minute is only 1s since the start", ErrNoSample)} + cp := &fakeCP{} + rec := runFor(t, 2*time.Minute, t.TempDir(), cp, collector) + + if n := len(rec.times("skipping the sample of this minute")); n != 0 { + t.Errorf("%d warnings, want none for a minute without a sample", n) + } + if got := rec.recordsOf("no sample for this minute"); len(got) != 2 || got[0].Level != slog.LevelInfo { + t.Errorf("%d info lines, want 2", len(got)) + } + // The reports still go: the events are sent. + if len(cp.eventsAt) == 0 { + t.Error("no events request, want the report to run") + } + }) +} + func TestTheReportTimeRunningOutIsNotAFailure(t *testing.T) { synctest.Test(t, func(t *testing.T) { dir := t.TempDir() diff --git a/internal/agent/wire/wire.go b/internal/agent/wire/wire.go index 78c1f2c..32e4673 100644 --- a/internal/agent/wire/wire.go +++ b/internal/agent/wire/wire.go @@ -45,6 +45,14 @@ type Sample struct { NetInBytes uint64 `json:"net_in_bytes"` NetOutBytes uint64 `json:"net_out_bytes"` NetCountersReset bool `json:"net_counters_reset"` + + // The peaks within the minute, from the readings each 10 seconds + // (contract v0.4.0). nil (JSON null) means "not known". + CPUMaxPercent *float64 `json:"cpu_max_percent"` + MemoryUsedMaxBytes *uint64 `json:"memory_used_max_bytes"` + 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"` } // MetricsReply is the reply to POST /agent/v1/metrics. diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index d87cb6f..1638365 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -1,20 +1,24 @@ // Package metrics measures a Linux server for the monitoring agent: CPU, -// load, memory, swap, disk and network each minute, and the status of the -// server. It needs no root. +// 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. package metrics import ( "bytes" "context" "errors" + "fmt" "io/fs" "log/slog" "os" "os/exec" "path/filepath" "runtime" + "slices" "time" + "github.com/flywp/server-cli/internal/agent" "github.com/flywp/server-cli/internal/agent/wire" "github.com/flywp/server-cli/internal/statefile" ) @@ -29,15 +33,26 @@ const ( updatesEvery = time.Hour aptCheckPath = "/usr/lib/update-notifier/apt-check" aptCheckTimeout = 30 * time.Second + + // minFirstMinute is the shortest first minute after a start. A shorter + // one is mostly the load of the start, for example an install: its CPU + // value would be a false spike. + minFirstMinute = 10 * time.Second + + // maxReadings limits the readings between two ticks. The agent reads + // five times between two ticks; more readings come only when ticks fail. + maxReadings = 30 ) -// counters is the previous reading. It is saved, so that the first sample -// after an agent restart continues from it. -type counters struct { +// reading holds the counters at one moment. The reading of the tick is saved, +// so that the first sample after an agent restart continues from it. +type reading struct { BootID string `json:"boot_id"` At time.Time `json:"at"` CPU cpuTimes `json:"cpu"` Net map[string]netCounters `json:"net"` + // mem is not saved: only the readings of the minute give its peak. + mem memory } // Collector measures the server. Use New. @@ -52,7 +67,15 @@ type Collector struct { release func() string aptCheck func(ctx context.Context) ([]byte, error) - prev *counters + // prev is the reading of the last tick, and readings are the readings + // after it, oldest first. They give the windows of the next sample. + prev *reading + readings []reading + // fromStart is true while prev is the reading of the start of this + // process, not a saved reading. minFirst is the shortest first minute; + // tests set it to 0. + fromStart bool + minFirst time.Duration // noInterface is true after the warning that no interface is counted. noInterface bool @@ -72,51 +95,95 @@ func New(root, stateDir string, log *slog.Logger) *Collector { statfs: statfs, release: kernelRelease, aptCheck: runAptCheck, + minFirst: minFirstMinute, } // Take a reading now, so that the first sample has a CPU value for the - // time since the start. The saved reading replaces it only when it is + // time since the start. The saved reading comes before it only when it is // recent and from this boot: then the traffic continues without a gap. now := time.Now() - if cur, err := c.read(now); err == nil { - cur.Net = nil - c.prev = &cur - } + start, startErr := c.read(now) - var saved counters + var saved reading err := statefile.Read(c.path, &saved) - switch { - case err == nil: - if c.prev != nil && saved.BootID == c.prev.BootID && now.Sub(saved.At) <= maxAge { - c.prev = &saved - } - case !errors.Is(err, fs.ErrNotExist): + if err != nil && !errors.Is(err, fs.ErrNotExist) { log.Warn("ignoring the saved counters", "error", err) } + switch { + case startErr != nil: + // The first sample has no previous reading. + case err == nil && saved.BootID == start.BootID && saved.At.Before(now) && now.Sub(saved.At) <= maxAge: + c.prev = &saved + c.readings = []reading{start} + default: + // The traffic since the start is not the traffic of one minute. + start.Net = nil + c.prev = &start + c.fromStart = true + } + return c } +// Read takes a reading between two ticks, for the peaks of the minute. A +// reading that fails is left out: the windows on each side of it join into +// one (contract v0.4.0). +func (c *Collector) Read(now time.Time) { + if last := c.last(); last != nil && !now.After(last.At) { + return + } + + r, err := c.read(now) + if err != nil { + c.log.Debug("skipping a reading", "error", err) + return + } + if len(c.readings) >= maxReadings { + c.readings = slices.Delete(c.readings, 0, 1) + } + c.readings = append(c.readings, r) +} + +// last returns the newest reading, or nil. +func (c *Collector) last() *reading { + if n := len(c.readings); n > 0 { + return &c.readings[n-1] + } + return c.prev +} + // Sample measures the minute that ends at now. At the first sample and then // each hour, it also counts the waiting updates for Status: apt-check takes -// some seconds, and a sample comes before the sends of a report. +// some seconds, and a sample comes before the sends of a report. The count +// comes after the reading of the tick, so that it does not move the reading. func (c *Collector) Sample(now time.Time) (wire.Sample, error) { - c.refreshUpdates(context.Background()) - cur, err := c.read(now) if err != nil { return wire.Sample{}, err } + // A tick right after the start has no minute to measure. Its reading + // starts the next minute, which then has all its values. + if c.fromStart { + c.fromStart = false + if d := cur.At.Sub(c.prev.At); d >= 0 && d < c.minFirst { + c.prev, c.readings = &cur, nil + if err := statefile.Write(c.path, cur); err != nil { + c.log.Warn("saving the counters", "error", err) + } + return wire.Sample{}, fmt.Errorf("%w: the first minute is only %s since the start", agent.ErrNoSample, d.Round(time.Millisecond)) + } + } + + c.refreshUpdates(context.Background()) + var s wire.Sample if s.Load1, err = parseFile(c, "proc/loadavg", parseLoad); err != nil { return wire.Sample{}, err } - mem, err := parseFile(c, "proc/meminfo", parseMeminfo) - if err != nil { - return wire.Sample{}, err - } + mem := cur.mem s.MemoryTotalBytes = mem.total s.MemoryUsedBytes = mem.total - min(mem.available, mem.total) s.SwapTotalBytes = mem.swapTotal @@ -140,7 +207,10 @@ func (c *Collector) Sample(now time.Time) (wire.Sample, error) { } } + setPeaks(&s, c.prev, c.readings, cur) + c.prev = &cur + c.readings = nil if err := statefile.Write(c.path, cur); err != nil { c.log.Warn("saving the counters", "error", err) } @@ -152,7 +222,7 @@ func (c *Collector) Sample(now time.Time) (wire.Sample, error) { // that both readings have, so a new or a removed interface makes no spike. // reset is true when the traffic of the minute is not known: no previous // reading, a reboot, a reading older than maxAge, or a counter that went back. -func netDelta(prev *counters, cur counters) (in, out uint64, reset bool) { +func netDelta(prev *reading, cur reading) (in, out uint64, reset bool) { if prev == nil || prev.Net == nil || prev.BootID != cur.BootID { return 0, 0, true } @@ -176,15 +246,20 @@ func netDelta(prev *counters, cur counters) (in, out uint64, reset bool) { } // read takes the counters now. -func (c *Collector) read(now time.Time) (counters, error) { +func (c *Collector) read(now time.Time) (reading, error) { cpu, err := parseFile(c, "proc/stat", parseCPU) if err != nil { - return counters{}, err + return reading{}, err + } + + mem, err := parseFile(c, "proc/meminfo", parseMeminfo) + if err != nil { + return reading{}, err } all, err := parseFile(c, "proc/net/dev", parseNetDev) if err != nil { - return counters{}, err + return reading{}, err } net := map[string]netCounters{} @@ -194,10 +269,10 @@ func (c *Collector) read(now time.Time) (counters, error) { bootID, err := os.ReadFile(c.file("proc/sys/kernel/random/boot_id")) if err != nil { - return counters{}, err + return reading{}, err } - return counters{BootID: string(bytes.TrimSpace(bootID)), At: now, CPU: cpu, Net: net}, nil + return reading{BootID: string(bytes.TrimSpace(bootID)), At: now, CPU: cpu, Net: net, mem: mem}, nil } // interfaces returns the network interfaces that have a hardware device and diff --git a/internal/metrics/metrics_linux_test.go b/internal/metrics/metrics_linux_test.go index 7837b59..8fd5bb1 100644 --- a/internal/metrics/metrics_linux_test.go +++ b/internal/metrics/metrics_linux_test.go @@ -11,6 +11,8 @@ import ( // TestRealServer measures this Linux machine: CI runs it on ubuntu-latest. func TestRealServer(t *testing.T) { c := New("/", t.TempDir(), slog.New(slog.DiscardHandler)) + // The sample comes right after the start: measure it anyway. + c.minFirst = 0 s, err := c.Sample(time.Now()) if err != nil { diff --git a/internal/metrics/metrics_test.go b/internal/metrics/metrics_test.go index c69630b..41163cc 100644 --- a/internal/metrics/metrics_test.go +++ b/internal/metrics/metrics_test.go @@ -13,6 +13,7 @@ import ( "testing" "time" + "github.com/flywp/server-cli/internal/agent" "github.com/flywp/server-cli/internal/agent/wire" ) @@ -83,6 +84,8 @@ func (s *server) collector(stateDir string) *Collector { c.statfs = func(string) (uint64, uint64, error) { return 100 << 30, 25 << 30, nil } c.release = func() string { return "6.8.0-45-generic" } c.aptCheck = func(context.Context) ([]byte, error) { return []byte("33;6"), nil } + // The tests take the first sample some milliseconds after the start. + c.minFirst = 0 return c } @@ -419,3 +422,63 @@ func TestNewWithCountersThatCannotBeRead(t *testing.T) { t.Error("net_counters_reset = false, want true after counters that cannot be read") } } + +func TestAShortFirstMinuteHasNoSample(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + c.minFirst = minFirstMinute + + // The install ends 1 s before the tick: that second is all busy. + srv.cpu(1100, 800) + first := time.Now().Add(time.Second) + if _, err := c.Sample(first); !errors.Is(err, agent.ErrNoSample) { + t.Fatalf("Sample() 1 s after the start = %v, want ErrNoSample", err) + } + + // The next minute is a whole minute from the tick reading, with its + // traffic. + srv.cpu(1700, 1250) + srv.net(map[string][2]uint64{"eth0": {1600, 800}, "eth1": {100, 50}}) + s, err := c.Sample(first.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.CPUPercent != 25 { + t.Errorf("cpu_percent = %v, want 25 from the tick reading, not from the start", s.CPUPercent) + } + if s.NetCountersReset || s.NetInBytes != 600 || s.NetInMaxBytesPerSecond == nil { + t.Errorf("net = %d, peak %v, reset %v; want 600 and a peak: the minute has a previous reading", s.NetInBytes, ptr(s.NetInMaxBytesPerSecond), s.NetCountersReset) + } +} + +func TestALongFirstMinuteHasASample(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + c.minFirst = minFirstMinute + + srv.cpu(1600, 950) + s, err := c.Sample(time.Now().Add(20 * time.Second)) + if err != nil { + t.Fatalf("Sample() 20 s after the start = %v, want a sample", err) + } + if s.CPUPercent != 75 { + t.Errorf("cpu_percent = %v, want 75", s.CPUPercent) + } +} + +func TestARestartWithSavedCountersHasASample(t *testing.T) { + srv := newServer(t) + state := t.TempDir() + now := time.Now() + if _, err := srv.collector(state).Sample(now.Add(-59 * time.Second)); err != nil { + t.Fatal(err) + } + + // A quick restart, for example an update, 1 s before the tick: the + // minute continues from the saved reading. + c := srv.collector(state) + c.minFirst = minFirstMinute + if _, err := c.Sample(now.Add(time.Second)); err != nil { + t.Errorf("Sample() after a restart with saved counters = %v, want a sample", err) + } +} diff --git a/internal/metrics/windows.go b/internal/metrics/windows.go new file mode 100644 index 0000000..5938fb9 --- /dev/null +++ b/internal/metrics/windows.go @@ -0,0 +1,108 @@ +package metrics + +import ( + "time" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +// minWindow is the shortest window between two readings. A reading closer +// than this to its neighbours is left out, because a peak of some +// milliseconds is noise, not the peak of 10 seconds. Only the minute itself +// can be shorter, for example when the agent started just before the tick. +const minWindow = 5 * time.Second + +// chain returns the readings that make the windows of the minute that ends +// at cur: prev (when it is before cur), the readings between, and cur, in +// time order. The times always go up, also after a step of the wall clock. +func chain(prev *reading, readings []reading, cur reading) []reading { + var out []reading + if prev != nil && prev.At.Before(cur.At) { + out = append(out, *prev) + } + for _, r := range readings { + if n := len(out); n > 0 && r.At.Sub(out[n-1].At) < minWindow { + continue + } + if cur.At.Sub(r.At) < minWindow { + continue + } + out = append(out, r) + } + + 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 +// 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) + + // The memory peaks come from the readings of the minute: not from the + // reading of the tick before, and not from an older minute whose tick + // failed. + var memMax, swapMax uint64 + for _, r := range all { + if cur.At.Sub(r.At) >= time.Minute { + continue + } + memMax = max(memMax, r.mem.total-min(r.mem.available, r.mem.total)) + swapMax = max(swapMax, r.mem.swapTotal-min(r.mem.swapFree, r.mem.swapTotal)) + } + memMax = min(max(memMax, s.MemoryUsedBytes), s.MemoryTotalBytes) + swapMax = min(max(swapMax, s.SwapUsedBytes), s.SwapTotalBytes) + s.MemoryUsedMaxBytes, s.SwapUsedMaxBytes = &memMax, &swapMax + + if cpu, ok := cpuPeak(all); ok { + cpu = max(cpu, s.CPUPercent) + s.CPUMaxPercent = &cpu + } + + if !s.NetCountersReset { + if in, out, ok := netPeaks(all); ok { + // The control plane reads the value of the minute as bytes / 60. + in, out = max(in, s.NetInBytes/60), max(out, s.NetOutBytes/60) + s.NetInMaxBytesPerSecond, s.NetOutMaxBytesPerSecond = &in, &out + } + } +} + +// cpuPeak returns the busy share of the busiest window. A window across a +// reboot, with counters that went back, or of an older minute whose tick +// failed, is left out. +func cpuPeak(all []reading) (peak float64, ok bool) { + cur := all[len(all)-1] + for i := 1; i < len(all); i++ { + a, b := all[i-1], all[i] + if cur.At.Sub(b.At) >= time.Minute { + continue + } + if a.BootID != b.BootID || b.CPU.Total <= a.CPU.Total || b.CPU.Idle < a.CPU.Idle { + continue + } + peak, ok = max(peak, cpuPercent(a.CPU, b.CPU)), true + } + + return peak, ok +} + +// netPeaks returns the received and sent bytes each second of the busiest +// windows, rounded down. It returns false when the traffic of a window is not +// known, for example after a counter went back. +func netPeaks(all []reading) (in, out uint64, ok bool) { + for i := 1; i < len(all); i++ { + a, b := all[i-1], all[i] + dIn, dOut, reset := netDelta(&a, b) + if reset { + return 0, 0, false + } + secs := b.At.Sub(a.At).Seconds() + in = max(in, uint64(float64(dIn)/secs)) + out = max(out, uint64(float64(dOut)/secs)) + ok = true + } + + return in, out, ok +} diff --git a/internal/metrics/windows_test.go b/internal/metrics/windows_test.go new file mode 100644 index 0000000..32f9e7d --- /dev/null +++ b/internal/metrics/windows_test.go @@ -0,0 +1,341 @@ +package metrics + +import ( + "fmt" + "os" + "path/filepath" + "testing" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +// mem writes /proc/meminfo with the available memory and the free swap, in +// kB. The totals are 8000000 kB and 2000000 kB. +func (s *server) mem(availableKB, swapFreeKB uint64) { + s.write("proc/meminfo", fmt.Sprintf("MemTotal: 8000000 kB\nMemFree: 500000 kB\nMemAvailable: %d kB\nSwapTotal: 2000000 kB\nSwapFree: %d kB\n", availableKB, swapFreeKB)) +} + +// step is one step of a test minute: the counters that the server shows at +// a reading. +type step struct { + total, idle uint64 // CPU ticks + in, out uint64 // eth0 bytes + available uint64 // kB + swapFree uint64 // kB +} + +func (s *server) set(st step) { + s.cpu(st.total, st.idle) + s.net(map[string][2]uint64{"eth0": {st.in, st.out}, "eth1": {100, 50}}) + s.mem(st.available, st.swapFree) +} + +// runMinute takes a sample at base, a reading each 10 seconds, and the +// sample of the next tick at base + 60 s. steps holds the counters of the +// five readings and of the tick. +func runMinute(t *testing.T, srv *server, c *Collector, base time.Time, first step, steps [6]step) wire.Sample { + t.Helper() + + srv.set(first) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + for i := range 5 { + srv.set(steps[i]) + c.Read(base.Add(time.Duration(i+1) * 10 * time.Second)) + } + srv.set(steps[5]) + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + return s +} + +func TestPeaks(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + // Each window adds 100 CPU ticks. The window from 20 s to 30 s is 90% + // busy; the others are 10% busy. eth0 receives 10000 bytes in the + // window from 40 s to 50 s, and 1000 bytes in each other window. The + // memory peaks at 30 s, and the swap at 40 s. + first := step{1000, 800, 0, 0, 6000000, 1500000} + s := runMinute(t, srv, c, base, first, [6]step{ + {1100, 890, 1000, 100, 6000000, 1500000}, + {1200, 980, 2000, 200, 5000000, 1500000}, + {1300, 990, 3000, 300, 4000000, 1500000}, + {1400, 1080, 4000, 400, 5000000, 1000000}, + {1500, 1170, 14000, 500, 6000000, 1500000}, + {1600, 1260, 15000, 600, 6000000, 1500000}, + }) + + // The minute: 600 ticks, 460 idle. + if want := float64(140) / 600 * 100; s.CPUPercent != want { + t.Errorf("cpu_percent = %v, want %v", s.CPUPercent, want) + } + if s.CPUMaxPercent == nil || *s.CPUMaxPercent != 90 { + t.Errorf("cpu_max_percent = %v, want 90", ptr(s.CPUMaxPercent)) + } + if want := uint64(4000000 * 1024); s.MemoryUsedMaxBytes == nil || *s.MemoryUsedMaxBytes != want { + t.Errorf("memory_used_max_bytes = %v, want %d", ptr(s.MemoryUsedMaxBytes), want) + } + if want := uint64(1000000 * 1024); s.SwapUsedMaxBytes == nil || *s.SwapUsedMaxBytes != want { + t.Errorf("swap_used_max_bytes = %v, want %d", ptr(s.SwapUsedMaxBytes), want) + } + if s.NetInMaxBytesPerSecond == nil || *s.NetInMaxBytesPerSecond != 1000 { + t.Errorf("net_in_max_bytes_per_second = %v, want 1000", ptr(s.NetInMaxBytesPerSecond)) + } + if s.NetOutMaxBytesPerSecond == nil || *s.NetOutMaxBytesPerSecond != 10 { + t.Errorf("net_out_max_bytes_per_second = %v, want 10", ptr(s.NetOutMaxBytesPerSecond)) + } + checkOrder(t, s) +} + +func TestAFailedReadingJoinsTwoWindows(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + srv.set(step{1000, 800, 0, 0, 6000000, 1500000}) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + + // The reading at 10 s fails. The window from 0 s to 20 s has 190 idle + // ticks of 200. + stat := filepath.Join(srv.root, "proc/stat") + if err := os.Remove(stat); err != nil { + t.Fatal(err) + } + c.Read(base.Add(10 * time.Second)) + srv.set(step{1200, 990, 0, 0, 6000000, 1500000}) + c.Read(base.Add(20 * time.Second)) + srv.set(step{1300, 1000, 0, 0, 6000000, 1500000}) + + s, err := c.Sample(base.Add(30 * time.Second)) + if err != nil { + t.Fatal(err) + } + if s.CPUMaxPercent == nil || *s.CPUMaxPercent != 90 { + t.Errorf("cpu_max_percent = %v, want 90 from the window of 20 s to 30 s", ptr(s.CPUMaxPercent)) + } + if s.NetInMaxBytesPerSecond == nil || *s.NetInMaxBytesPerSecond != 0 { + t.Errorf("net_in_max_bytes_per_second = %v, want 0", ptr(s.NetInMaxBytesPerSecond)) + } +} + +func TestFirstSamplePeaksEqualTheMinute(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + + // The agent started just now: one window, from the start to the tick. + srv.set(step{1600, 950, 5000, 5000, 5000000, 1500000}) + s, err := c.Sample(time.Now().Add(20 * time.Second)) + if err != nil { + t.Fatal(err) + } + if s.CPUMaxPercent == nil || *s.CPUMaxPercent != s.CPUPercent || s.CPUPercent != 75 { + t.Errorf("cpu = %v, max %v; want 75 and an equal peak", s.CPUPercent, ptr(s.CPUMaxPercent)) + } + if s.MemoryUsedMaxBytes == nil || *s.MemoryUsedMaxBytes != s.MemoryUsedBytes { + t.Errorf("memory_used_max_bytes = %v, want the value of the minute %d", ptr(s.MemoryUsedMaxBytes), s.MemoryUsedBytes) + } + // The traffic since the start is not known, so its peak is not known. + if !s.NetCountersReset || s.NetInMaxBytesPerSecond != nil || s.NetOutMaxBytesPerSecond != nil { + t.Errorf("net peaks = %v, %v with reset %v; want null", ptr(s.NetInMaxBytesPerSecond), ptr(s.NetOutMaxBytesPerSecond), s.NetCountersReset) + } +} + +func TestPeaksAfterARestartIncludeTheSavedReading(t *testing.T) { + srv := newServer(t) + state := t.TempDir() + now := time.Now() + + // The last tick of the old process was 30 s ago. + srv.set(step{1000, 800, 0, 0, 6000000, 1500000}) + if _, err := srv.collector(state).Sample(now.Add(-30 * time.Second)); err != nil { + t.Fatal(err) + } + + // The new process starts after 100% busy time. Then it takes the tick. + srv.set(step{1100, 800, 3000, 0, 6000000, 1500000}) + c := srv.collector(state) + srv.set(step{1600, 1250, 3000, 0, 6000000, 1500000}) + s, err := c.Sample(now.Add(30 * time.Second)) + if err != nil { + t.Fatal(err) + } + if s.CPUPercent != 25 { + t.Errorf("cpu_percent = %v, want 25 from the saved reading", s.CPUPercent) + } + if s.CPUMaxPercent == nil || *s.CPUMaxPercent != 100 { + t.Errorf("cpu_max_percent = %v, want 100 from the window of the saved reading to the start", ptr(s.CPUMaxPercent)) + } + if s.NetCountersReset || s.NetInBytes != 3000 || s.NetInMaxBytesPerSecond == nil || *s.NetInMaxBytesPerSecond < 99 { + t.Errorf("net = %d, peak %v, reset %v; want 3000 with a peak of about 100 each second (3000 bytes in 30 s)", s.NetInBytes, ptr(s.NetInMaxBytesPerSecond), s.NetCountersReset) + } + checkOrder(t, s) +} + +func TestNetPeaksAreNullAfterACounterWentBack(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + // eth0 goes back at 30 s and forward again by the tick: the minute + // looks correct, but a window is not known. + s := runMinute(t, srv, c, base, step{1000, 800, 5000, 5000, 6000000, 1500000}, [6]step{ + {1100, 890, 6000, 6000, 6000000, 1500000}, + {1200, 980, 7000, 7000, 6000000, 1500000}, + {1300, 990, 10, 10, 6000000, 1500000}, + {1400, 1080, 1000, 1000, 6000000, 1500000}, + {1500, 1170, 6000, 6000, 6000000, 1500000}, + {1600, 1260, 9000, 9000, 6000000, 1500000}, + }) + if s.NetInMaxBytesPerSecond != nil || s.NetOutMaxBytesPerSecond != nil { + t.Errorf("net peaks = %v, %v; want null after a window with a counter that went back", ptr(s.NetInMaxBytesPerSecond), ptr(s.NetOutMaxBytesPerSecond)) + } + if s.CPUMaxPercent == nil { + t.Error("cpu_max_percent = null, want a value: only the network is not known") + } +} + +func TestReadingsThatAreNotNewerAreIgnored(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + srv.set(step{1000, 800, 0, 0, 6000000, 1500000}) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + c.Read(base.Add(20 * time.Second)) + // The wall clock stepped back: the same reading comes again, and an + // older one. + c.Read(base.Add(20 * time.Second)) + c.Read(base.Add(10 * time.Second)) + if len(c.readings) != 1 { + t.Errorf("%d readings, want 1", len(c.readings)) + } + + // A reading close to the tick makes no window of some milliseconds. + c.Read(base.Add(59*time.Second + 900*time.Millisecond)) + srv.set(step{1100, 810, 0, 0, 6000000, 1500000}) + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.CPUMaxPercent == nil || *s.CPUMaxPercent != 90 { + t.Errorf("cpu_max_percent = %v, want 90 from the window of 20 s to 60 s", ptr(s.CPUMaxPercent)) + } +} + +func TestMemoryPeakAfterAFailedTickIsFromTheLastMinute(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + srv.set(step{1000, 800, 0, 0, 6000000, 1500000}) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + // A high use of memory at 30 s. The tick at 60 s fails. + srv.mem(1000000, 1500000) + c.Read(base.Add(30 * time.Second)) + srv.mem(6000000, 1500000) + c.Read(base.Add(50 * time.Second)) + stat := filepath.Join(srv.root, "proc/stat") + if err := os.Remove(stat); err != nil { + t.Fatal(err) + } + if _, err := c.Sample(base.Add(time.Minute)); err == nil { + t.Fatal("Sample() = nil error, want the error of the tick") + } + + c.Read(base.Add(70 * time.Second)) + srv.set(step{1100, 900, 0, 0, 6000000, 1500000}) + s, err := c.Sample(base.Add(2 * time.Minute)) + if err != nil { + t.Fatal(err) + } + if want := uint64(2000000 * 1024); s.MemoryUsedMaxBytes == nil || *s.MemoryUsedMaxBytes != want { + t.Errorf("memory_used_max_bytes = %v, want %d: the peak at 30 s is in the minute before", ptr(s.MemoryUsedMaxBytes), want) + } +} + +func TestCPUPeakAfterAFailedTickIsFromTheLastMinute(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + + srv.set(step{1000, 800, 0, 0, 6000000, 1500000}) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + // 100% busy from 0 s to 10 s. The tick at 60 s fails. + srv.set(step{1100, 800, 0, 0, 6000000, 1500000}) + c.Read(base.Add(10 * time.Second)) + if err := os.Remove(filepath.Join(srv.root, "proc/meminfo")); err != nil { + t.Fatal(err) + } + if _, err := c.Sample(base.Add(time.Minute)); err == nil { + t.Fatal("Sample() = nil error, want the error of the tick") + } + + // The next minute is idle. + srv.set(step{1100, 800, 0, 0, 6000000, 1500000}) + c.Read(base.Add(70 * time.Second)) + srv.set(step{1200, 900, 0, 0, 6000000, 1500000}) + s, err := c.Sample(base.Add(2 * time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.CPUMaxPercent == nil || *s.CPUMaxPercent >= 100 { + t.Errorf("cpu_max_percent = %v, want the peak of the last minute, not the 100%% of the minute before", ptr(s.CPUMaxPercent)) + } + checkOrder(t, s) +} + +func TestReadingsAreLimited(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + base := time.Now().Add(time.Second) + for i := range 100 { + c.Read(base.Add(time.Duration(i) * 10 * time.Second)) + } + if len(c.readings) != maxReadings { + t.Errorf("%d readings, want at most %d", len(c.readings), maxReadings) + } +} + +// checkOrder checks the order that the contract guarantees between a value of +// the minute and its peak. +func checkOrder(t *testing.T, s wire.Sample) { + t.Helper() + if s.CPUMaxPercent != nil && s.CPUPercent > *s.CPUMaxPercent { + t.Errorf("cpu_percent %v > cpu_max_percent %v", s.CPUPercent, *s.CPUMaxPercent) + } + if m := s.MemoryUsedMaxBytes; m == nil || s.MemoryUsedBytes > *m || *m > s.MemoryTotalBytes { + t.Errorf("memory %d, max %v, total %d: want used ≤ max ≤ total", s.MemoryUsedBytes, ptr(m), s.MemoryTotalBytes) + } + if m := s.SwapUsedMaxBytes; m == nil || s.SwapUsedBytes > *m || *m > s.SwapTotalBytes { + t.Errorf("swap %d, max %v, total %d: want used ≤ max ≤ total", s.SwapUsedBytes, ptr(m), s.SwapTotalBytes) + } + if m := s.NetInMaxBytesPerSecond; m != nil && s.NetInBytes > 60**m+59 { + t.Errorf("net_in_bytes %d > 60 × %d", s.NetInBytes, *m) + } + if m := s.NetOutMaxBytesPerSecond; m != nil && s.NetOutBytes > 60**m+59 { + t.Errorf("net_out_bytes %d > 60 × %d", s.NetOutBytes, *m) + } +} + +// ptr shows a value that can be null. +func ptr[T any](v *T) any { + if v == nil { + return "null" + } + return *v +}