diff --git a/cmd/agent.go b/cmd/agent.go index b7f4117..0afe1ba 100644 --- a/cmd/agent.go +++ b/cmd/agent.go @@ -7,6 +7,7 @@ import ( "syscall" "github.com/flywp/server-cli/internal/agent" + "github.com/flywp/server-cli/internal/metrics" "github.com/spf13/cobra" ) @@ -33,7 +34,8 @@ FLY_AGENT_SERVER_ID and STATE_DIRECTORY from the environment.`, ctx, stop := signal.NotifyContext(cmd.Context(), syscall.SIGTERM, os.Interrupt) defer stop() - return agent.Run(ctx, cfg, slog.New(slog.NewTextHandler(os.Stderr, nil)), nil) + log := slog.New(slog.NewTextHandler(os.Stderr, nil)) + return agent.Run(ctx, cfg, log, metrics.New("/", cfg.StateDir, log)) }, } diff --git a/go.mod b/go.mod index 8cf103c..cd772e0 100644 --- a/go.mod +++ b/go.mod @@ -9,6 +9,7 @@ require ( github.com/oklog/ulid/v2 v2.1.2 github.com/spf13/cobra v1.10.2 golang.org/x/mod v0.41.0 + golang.org/x/sys v0.48.0 gopkg.in/yaml.v2 v2.4.0 ) @@ -17,5 +18,4 @@ require ( github.com/mattn/go-colorable v0.1.15 // indirect github.com/mattn/go-isatty v0.0.24 // indirect github.com/spf13/pflag v1.0.10 // indirect - golang.org/x/sys v0.48.0 // indirect ) diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go new file mode 100644 index 0000000..57710b6 --- /dev/null +++ b/internal/metrics/metrics.go @@ -0,0 +1,322 @@ +// 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. +package metrics + +import ( + "bytes" + "context" + "errors" + "io/fs" + "log/slog" + "os" + "os/exec" + "path/filepath" + "runtime" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" + "github.com/flywp/server-cli/internal/statefile" +) + +const ( + // maxAge is the oldest previous reading that gives the traffic of one + // minute. The agent samples each 60 seconds; an older reading covers + // more than one minute. + maxAge = 90 * time.Second + + // The update counts come from apt-check, which takes some seconds. + updatesEvery = time.Hour + aptCheckPath = "/usr/lib/update-notifier/apt-check" + aptCheckTimeout = 30 * time.Second +) + +// counters is the previous reading. It is saved, so that the first sample +// after an agent restart continues from it. +type counters struct { + BootID string `json:"boot_id"` + At time.Time `json:"at"` + CPU cpuTimes `json:"cpu"` + Net map[string]netCounters `json:"net"` +} + +// Collector measures the server. Use New. +type Collector struct { + root string + path string // counters.json + log *slog.Logger + + // statfs returns the size and the used space of the file system of a + // path, and release the kernel release. Tests replace them. + statfs func(path string) (total, used uint64, err error) + release func() string + aptCheck func(ctx context.Context) ([]byte, error) + + prev *counters + // noInterface is true after the warning that no interface is counted. + noInterface bool + + updatesAt time.Time + updatesKnown bool + updatesTotal uint64 + updatesSecurity uint64 +} + +// New returns a collector that reads the files under root ("/" on a server) +// and keeps its counters in stateDir. +func New(root, stateDir string, log *slog.Logger) *Collector { + c := &Collector{ + root: root, + path: filepath.Join(stateDir, "counters.json"), + log: log, + statfs: statfs, + release: kernelRelease, + aptCheck: runAptCheck, + } + + // 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 + // 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 + } + + var saved counters + 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): + log.Warn("ignoring the saved counters", "error", err) + } + + return c +} + +// 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. +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 + } + + 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 + } + s.MemoryTotalBytes = mem.total + s.MemoryUsedBytes = mem.total - min(mem.available, mem.total) + s.SwapTotalBytes = mem.swapTotal + s.SwapUsedBytes = mem.swapTotal - min(mem.swapFree, mem.swapTotal) + + if s.DiskTotalBytes, s.DiskUsedBytes, err = c.statfs(c.file("")); err != nil { + return wire.Sample{}, err + } + + prev := c.prev + if prev != nil && prev.BootID == cur.BootID && cur.CPU.Total >= prev.CPU.Total { + s.CPUPercent = cpuPercent(prev.CPU, cur.CPU) + } + s.NetInBytes, s.NetOutBytes, s.NetCountersReset = netDelta(prev, cur) + if len(cur.Net) == 0 { + // No interface is counted, so the traffic is not known: it is not 0. + s.NetCountersReset = true + if !c.noInterface { + c.noInterface = true + c.log.Warn("no network interface to count: no interface has a hardware device, and the default route has none") + } + } + + c.prev = &cur + if err := statefile.Write(c.path, cur); err != nil { + c.log.Warn("saving the counters", "error", err) + } + + return s, nil +} + +// netDelta returns the traffic between two readings. It adds the interfaces +// 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) { + if prev == nil || prev.Net == nil || prev.BootID != cur.BootID { + return 0, 0, true + } + if age := cur.At.Sub(prev.At); age <= 0 || age > maxAge { + return 0, 0, true + } + + for name, c := range cur.Net { + p, ok := prev.Net[name] + if !ok { + continue + } + if c.In < p.In || c.Out < p.Out { + return 0, 0, true + } + in += c.In - p.In + out += c.Out - p.Out + } + + return in, out, false +} + +// read takes the counters now. +func (c *Collector) read(now time.Time) (counters, error) { + cpu, err := parseFile(c, "proc/stat", parseCPU) + if err != nil { + return counters{}, err + } + + all, err := parseFile(c, "proc/net/dev", parseNetDev) + if err != nil { + return counters{}, err + } + + net := map[string]netCounters{} + for _, name := range c.interfaces(all) { + net[name] = all[name] + } + + bootID, err := os.ReadFile(c.file("proc/sys/kernel/random/boot_id")) + if err != nil { + return counters{}, err + } + + return counters{BootID: string(bytes.TrimSpace(bootID)), At: now, CPU: cpu, Net: net}, nil +} + +// interfaces returns the network interfaces that have a hardware device and +// are not a port of an other interface. Thus lo, docker0, the Docker bridges +// and the veth interfaces are left out, and container traffic is not counted +// two or three times. A port (of a bond or a bridge, or the Azure VF under +// its netvsc interface) is left out too, because its traffic is also in the +// interface above it. If no interface is left, it returns the interface of +// the default route, for example the bond or the bridge. +func (c *Collector) interfaces(all map[string]netCounters) []string { + var names []string + for name := range all { + if _, err := os.Lstat(c.file("sys/class/net", name, "device")); err != nil { + continue + } + if _, err := os.Lstat(c.file("sys/class/net", name, "master")); err == nil { + continue + } + names = append(names, name) + } + if len(names) > 0 { + return names + } + + route, err := os.ReadFile(c.file("proc/net/route")) + if err != nil { + return nil + } + if name := parseDefaultRoute(route); name != "" { + if _, ok := all[name]; ok { + return []string{name} + } + } + + return nil +} + +// Status describes the server now. A value that cannot be read stays empty +// or 0, and the problem goes to the log. The update counts come from the last +// Sample. +func (c *Collector) Status(context.Context) wire.Status { + s := wire.Status{Arch: runtime.GOARCH, Kernel: c.release()} + + if _, err := os.Stat(c.file("var/run/reboot-required")); err == nil { + s.RebootRequired = true + } + + if data, err := os.ReadFile(c.file("etc/os-release")); err == nil { + s.OS = parseOSRelease(data) + } else { + c.log.Warn("reading the OS name", "error", err) + } + + if up, err := parseFile(c, "proc/uptime", parseUptime); err == nil { + s.UptimeSeconds = up + } else { + c.log.Warn("reading the uptime", "error", err) + } + + s.UpdatesTotal, s.UpdatesSecurity = c.updatesTotal, c.updatesSecurity + + return s +} + +// refreshUpdates counts the waiting updates at the first call and then each +// hour. If apt-check fails, the last counts stay. Without any count, they are +// 0: the contract has no value for "not known" yet. +func (c *Collector) refreshUpdates(ctx context.Context) { + if !c.updatesAt.IsZero() && time.Since(c.updatesAt) < updatesEvery { + return + } + c.updatesAt = time.Now() + + ctx, cancel := context.WithTimeout(ctx, aptCheckTimeout) + defer cancel() + + out, err := c.aptCheck(ctx) + var total, security uint64 + if err == nil { + total, security, err = parseAptCheck(out) + } + if err != nil { + if c.updatesKnown { + c.log.Warn("cannot count the waiting updates; keeping the last counts", "error", err) + } else { + c.log.Warn("cannot count the waiting updates; sending 0", "error", err) + } + return + } + + c.updatesKnown = true + c.updatesTotal, c.updatesSecurity = total, security +} + +// runAptCheck runs apt-check. It writes its result to stderr. +func runAptCheck(ctx context.Context) ([]byte, error) { + var stderr bytes.Buffer + cmd := exec.CommandContext(ctx, aptCheckPath) + cmd.Stderr = &stderr + cmd.WaitDelay = time.Second + if err := cmd.Run(); err != nil { + return nil, err + } + + return stderr.Bytes(), nil +} + +// file returns the path of a file under the root. +func (c *Collector) file(parts ...string) string { + return filepath.Join(append([]string{c.root}, parts...)...) +} + +// parseFile reads a file under the root and parses it. +func parseFile[T any](c *Collector, name string, parse func([]byte) (T, error)) (T, error) { + data, err := os.ReadFile(c.file(name)) + if err != nil { + var zero T + return zero, err + } + + return parse(data) +} diff --git a/internal/metrics/metrics_linux_test.go b/internal/metrics/metrics_linux_test.go new file mode 100644 index 0000000..7837b59 --- /dev/null +++ b/internal/metrics/metrics_linux_test.go @@ -0,0 +1,34 @@ +package metrics + +import ( + "context" + "log/slog" + "runtime" + "testing" + "time" +) + +// TestRealServer measures this Linux machine: CI runs it on ubuntu-latest. +func TestRealServer(t *testing.T) { + c := New("/", t.TempDir(), slog.New(slog.DiscardHandler)) + + s, err := c.Sample(time.Now()) + if err != nil { + t.Fatal(err) + } + + if s.CPUPercent < 0 || s.CPUPercent > 100 { + t.Errorf("cpu_percent = %v, want 0 to 100", s.CPUPercent) + } + if s.MemoryTotalBytes == 0 || s.MemoryUsedBytes > s.MemoryTotalBytes { + t.Errorf("memory = %d of %d", s.MemoryUsedBytes, s.MemoryTotalBytes) + } + if s.DiskTotalBytes == 0 || s.DiskUsedBytes > s.DiskTotalBytes { + t.Errorf("disk = %d of %d", s.DiskUsedBytes, s.DiskTotalBytes) + } + + st := c.Status(context.Background()) + if st.Kernel == "" || st.Arch != runtime.GOARCH || st.UptimeSeconds == 0 { + t.Errorf("status = %+v, want the kernel, the arch and the uptime", st) + } +} diff --git a/internal/metrics/metrics_test.go b/internal/metrics/metrics_test.go new file mode 100644 index 0000000..726fc2d --- /dev/null +++ b/internal/metrics/metrics_test.go @@ -0,0 +1,413 @@ +package metrics + +import ( + "context" + "errors" + "fmt" + "log/slog" + "os" + "path/filepath" + "runtime" + "sort" + "strings" + "testing" + "time" +) + +// server is a fake file system root with the files that the collector reads. +type server struct { + t *testing.T + root string +} + +func newServer(t *testing.T) *server { + t.Helper() + + s := &server{t: t, root: t.TempDir()} + s.write("proc/loadavg", "0.42 0.30 0.25 1/345 6789\n") + s.write("proc/meminfo", "MemTotal: 8000000 kB\nMemFree: 500000 kB\nMemAvailable: 6000000 kB\nSwapTotal: 2000000 kB\nSwapFree: 1500000 kB\n") + s.write("proc/uptime", "1892344.51 3700000.00\n") + s.write("proc/sys/kernel/random/boot_id", "boot-1\n") + s.write("proc/net/route", "Iface\tDestination\tGateway \tFlags\tRefCnt\tUse\tMetric\tMask\t\tMTU\tWindow\tIRTT\n"+ + "ens3\t00000000\t0101A8C0\t0003\t0\t0\t100\t00000000\t0\t0\t0\n"+ + "ens3\t0001A8C0\t00000000\t0001\t0\t0\t100\t00FFFFFF\t0\t0\t0\n") + s.write("etc/os-release", "NAME=\"Ubuntu\"\nPRETTY_NAME=\"Ubuntu 24.04.1 LTS\"\nID=ubuntu\n") + s.device("eth0") + s.device("eth1") + s.cpu(1000, 800) + s.net(map[string][2]uint64{"eth0": {1000, 500}, "eth1": {100, 50}, "lo": {9999, 9999}, "docker0": {7000, 7000}, "veth1": {7000, 7000}}) + return s +} + +func (s *server) write(name, content string) { + s.t.Helper() + path := filepath.Join(s.root, name) + if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { + s.t.Fatal(err) + } + if err := os.WriteFile(path, []byte(content), 0o644); err != nil { + s.t.Fatal(err) + } +} + +// device marks an interface as a hardware device. +func (s *server) device(name string) { + s.t.Helper() + if err := os.MkdirAll(filepath.Join(s.root, "sys/class/net", name, "device"), 0o755); err != nil { + s.t.Fatal(err) + } +} + +// cpu writes /proc/stat with the total and the idle ticks. +func (s *server) cpu(total, idle uint64) { + // user nice system idle iowait irq softirq steal guest guest_nice + busy := total - idle + s.write("proc/stat", fmt.Sprintf("cpu %d 0 0 %d 0 0 0 0 55 0\ncpu0 1 2 3 4 5 6 7 8 9 10\n", busy, idle)) +} + +// net writes /proc/net/dev with the received and sent bytes of each interface. +func (s *server) net(ifaces map[string][2]uint64) { + var b strings.Builder + b.WriteString("Inter-| Receive | Transmit\n") + b.WriteString(" face |bytes packets errs drop fifo frame compressed multicast|bytes packets errs drop fifo colls carrier compressed\n") + for name, c := range ifaces { + fmt.Fprintf(&b, "%6s: %d 10 0 0 0 0 0 0 %d 10 0 0 0 0 0 0\n", name, c[0], c[1]) + } + s.write("proc/net/dev", b.String()) +} + +func (s *server) collector(stateDir string) *Collector { + c := New(s.root, stateDir, slog.New(slog.DiscardHandler)) + 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 } + return c +} + +func TestSample(t *testing.T) { + srv := newServer(t) + now := time.Now() + c := srv.collector(t.TempDir()) + if _, err := c.Sample(now); err != nil { + t.Fatal(err) + } + + // One minute later: 600 more ticks with 150 idle, and some traffic. + srv.cpu(1600, 950) + srv.net(map[string][2]uint64{"eth0": {3000, 1500}, "eth1": {600, 150}, "lo": {99999, 99999}, "docker0": {70000, 70000}, "veth1": {70000, 70000}}) + + s, err := c.Sample(now.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + + if s.CPUPercent != 75 { + t.Errorf("cpu_percent = %v, want 75 (450 of 600 ticks busy)", s.CPUPercent) + } + if s.Load1 != 0.42 { + t.Errorf("load_1 = %v, want 0.42", s.Load1) + } + // MemTotal − MemAvailable, not MemTotal − MemFree. + if s.MemoryTotalBytes != 8000000*1024 || s.MemoryUsedBytes != 2000000*1024 { + t.Errorf("memory = %d of %d, want %d of %d", s.MemoryUsedBytes, s.MemoryTotalBytes, 2000000*1024, 8000000*1024) + } + if s.SwapTotalBytes != 2000000*1024 || s.SwapUsedBytes != 500000*1024 { + t.Errorf("swap = %d of %d, want %d of %d", s.SwapUsedBytes, s.SwapTotalBytes, 500000*1024, 2000000*1024) + } + if s.DiskTotalBytes != 100<<30 || s.DiskUsedBytes != 25<<30 { + t.Errorf("disk = %d of %d, want the statfs values", s.DiskUsedBytes, s.DiskTotalBytes) + } + // Only eth0 and eth1 have a device: lo, docker0 and veth1 are not counted. + if s.NetInBytes != 2000+500 || s.NetOutBytes != 1000+100 || s.NetCountersReset { + t.Errorf("net = in %d, out %d, reset %v; want in 2500, out 1100 from eth0 and eth1", s.NetInBytes, s.NetOutBytes, s.NetCountersReset) + } +} + +func TestFirstSampleHasNoTraffic(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + + s, err := c.Sample(time.Now()) + if err != nil { + t.Fatal(err) + } + if !s.NetCountersReset || s.NetInBytes != 0 || s.NetOutBytes != 0 { + t.Errorf("first sample net = %d, %d, reset %v; want 0, 0 and a reset", s.NetInBytes, s.NetOutBytes, s.NetCountersReset) + } +} + +func TestRestartContinuesFromTheSavedCounters(t *testing.T) { + srv := newServer(t) + state := t.TempDir() + now := time.Now() + if _, err := srv.collector(state).Sample(now); err != nil { + t.Fatal(err) + } + + // A new agent process starts, and one minute after the last sample it + // takes the next one. + srv.net(map[string][2]uint64{"eth0": {1500, 700}, "eth1": {100, 50}}) + s, err := srv.collector(state).Sample(now.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.NetCountersReset || s.NetInBytes != 500 || s.NetOutBytes != 200 { + t.Errorf("net after a restart = %d, %d, reset %v; want 500, 200 without a reset", s.NetInBytes, s.NetOutBytes, s.NetCountersReset) + } +} + +func TestTrafficResets(t *testing.T) { + tests := []struct { + name string + change func(*server) + after time.Duration + }{ + {"reboot", func(s *server) { s.write("proc/sys/kernel/random/boot_id", "boot-2\n") }, time.Minute}, + {"counter went back", func(s *server) { s.net(map[string][2]uint64{"eth0": {10, 10}, "eth1": {100, 50}}) }, 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) + now := time.Now() + c := srv.collector(t.TempDir()) + 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) + } + if !s.NetCountersReset || s.NetInBytes != 0 || s.NetOutBytes != 0 { + t.Errorf("net = %d, %d, reset %v; want 0, 0 and a reset", s.NetInBytes, s.NetOutBytes, s.NetCountersReset) + } + }) + } +} + +func TestStaleSavedCountersAreIgnored(t *testing.T) { + srv := newServer(t) + state := t.TempDir() + if _, err := srv.collector(state).Sample(time.Now().Add(-10 * time.Minute)); err != nil { + t.Fatal(err) + } + + // The agent was stopped for 10 minutes: its saved traffic counters do not + // give the traffic of one minute. + srv.net(map[string][2]uint64{"eth0": {900000, 900000}, "eth1": {100, 50}}) + s, err := srv.collector(state).Sample(time.Now()) + if err != nil { + t.Fatal(err) + } + if !s.NetCountersReset || s.NetInBytes != 0 { + t.Errorf("net = %d, reset %v; want 0 and a reset", s.NetInBytes, s.NetCountersReset) + } +} + +func TestNewInterfaceMakesNoSpike(t *testing.T) { + srv := newServer(t) + now := time.Now() + c := srv.collector(t.TempDir()) + if _, err := c.Sample(now); err != nil { + t.Fatal(err) + } + + srv.device("eth2") + srv.net(map[string][2]uint64{"eth0": {1100, 600}, "eth1": {100, 50}, "eth2": {5 << 40, 5 << 40}}) + s, err := c.Sample(now.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.NetInBytes != 100 || s.NetOutBytes != 100 { + t.Errorf("net = %d, %d; want 100, 100 (the new eth2 counts from the next sample)", s.NetInBytes, s.NetOutBytes) + } +} + +func TestInterfacesFallBackToTheDefaultRoute(t *testing.T) { + srv := newServer(t) + if err := os.RemoveAll(filepath.Join(srv.root, "sys")); err != nil { + t.Fatal(err) + } + srv.net(map[string][2]uint64{"ens3": {1, 1}, "lo": {1, 1}, "docker0": {1, 1}}) + + all, err := parseFile(srv.collector(t.TempDir()), "proc/net/dev", parseNetDev) + if err != nil { + t.Fatal(err) + } + if got := srv.collector(t.TempDir()).interfaces(all); fmt.Sprint(got) != "[ens3]" { + t.Errorf("interfaces() = %v, want [ens3] from the default route", got) + } +} + +func TestSampleErrors(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + if err := os.Remove(filepath.Join(srv.root, "proc/meminfo")); err != nil { + t.Fatal(err) + } + if _, err := c.Sample(time.Now()); err == nil { + t.Error("Sample() = nil error, want an error without /proc/meminfo") + } + + srv = newServer(t) + c = srv.collector(t.TempDir()) + c.statfs = func(string) (uint64, uint64, error) { return 0, 0, errors.New("statfs failed") } + if _, err := c.Sample(time.Now()); err == nil { + t.Error("Sample() = nil error, want the statfs error") + } +} + +func TestStatus(t *testing.T) { + srv := newServer(t) + srv.write("var/run/reboot-required", "*** System restart required ***\n") + c := srv.collector(t.TempDir()) + + calls := 0 + c.aptCheck = func(context.Context) ([]byte, error) { + calls++ + return []byte("33;6"), nil + } + + // The count runs with the first sample, before the sends of a report. + if _, err := c.Sample(time.Now()); err != nil { + t.Fatal(err) + } + s := c.Status(context.Background()) + if !s.RebootRequired || s.UpdatesTotal != 33 || s.UpdatesSecurity != 6 { + t.Errorf("status = %+v, want a restart and 33 updates with 6 security updates", s) + } + if s.OS != "Ubuntu 24.04.1 LTS" || s.Kernel != "6.8.0-45-generic" || s.UptimeSeconds != 1892344 || s.Arch != runtime.GOARCH { + t.Errorf("status = %+v", s) + } + + // The counts stay for one hour: apt-check takes some seconds. + c.Status(context.Background()) + if _, err := c.Sample(time.Now()); err != nil { + t.Fatal(err) + } + if calls != 1 { + t.Errorf("apt-check ran %d times, want 1 in one hour", calls) + } + + c.updatesAt = time.Now().Add(-2 * time.Hour) + if _, err := c.Sample(time.Now()); err != nil { + t.Fatal(err) + } + if calls != 2 { + t.Errorf("apt-check ran %d times, want again after one hour", calls) + } +} + +func TestStatusWithoutAptCheck(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + c.aptCheck = func(context.Context) ([]byte, error) { return nil, os.ErrNotExist } + + if _, err := c.Sample(time.Now()); err != nil { + t.Fatal(err) + } + s := c.Status(context.Background()) + if s.UpdatesTotal != 0 || s.UpdatesSecurity != 0 || s.RebootRequired { + t.Errorf("status = %+v, want 0 updates and no restart", s) + } + if s.OS == "" { + t.Error("status has no OS: one missing value must not clear the others") + } +} + +func TestAptCheckWithWarningsBeforeTheResult(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + // apt-check on Ubuntu 24.04 with a source that is configured two times. + c.aptCheck = func(context.Context) ([]byte, error) { + return []byte("/usr/lib/update-notifier/apt-check:351: Warning: W:Target Packages (main/binary-amd64/Packages) is configured multiple times in /etc/apt/sources.list:1 and /etc/apt/sources.list.d/ubuntu.sources:1\n" + + " apt_pkg.init()\n29;26"), nil + } + + if _, err := c.Sample(time.Now()); err != nil { + t.Fatal(err) + } + if s := c.Status(context.Background()); s.UpdatesTotal != 29 || s.UpdatesSecurity != 26 { + t.Errorf("updates = %d;%d, want 29;26 from the last line", s.UpdatesTotal, s.UpdatesSecurity) + } +} + +func TestAFailedCountKeepsTheLastCounts(t *testing.T) { + srv := newServer(t) + c := srv.collector(t.TempDir()) + if _, err := c.Sample(time.Now()); err != nil { + t.Fatal(err) + } + + // One hour later apt-check fails, for example with a timeout. + c.aptCheck = func(context.Context) ([]byte, error) { return nil, context.DeadlineExceeded } + c.updatesAt = time.Now().Add(-2 * time.Hour) + if _, err := c.Sample(time.Now()); err != nil { + t.Fatal(err) + } + if s := c.Status(context.Background()); s.UpdatesTotal != 33 || s.UpdatesSecurity != 6 { + t.Errorf("updates = %d;%d, want the last counts 33;6", s.UpdatesTotal, s.UpdatesSecurity) + } +} + +func TestNoCountedInterfaceIsNotZeroTraffic(t *testing.T) { + srv := newServer(t) + if err := os.RemoveAll(filepath.Join(srv.root, "sys")); err != nil { + t.Fatal(err) + } + // An IPv6-only host: no IPv4 default route. + srv.write("proc/net/route", "Iface\tDestination\tGateway \tFlags\tRefCnt\tUse\tMetric\tMask\t\tMTU\tWindow\tIRTT\n") + now := time.Now() + c := srv.collector(t.TempDir()) + if _, err := c.Sample(now); err != nil { + t.Fatal(err) + } + + s, err := c.Sample(now.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if !s.NetCountersReset { + t.Error("net_counters_reset = false, want true: the traffic is not known") + } +} + +func TestPortsOfAnOtherInterfaceAreNotCounted(t *testing.T) { + srv := newServer(t) + // Azure accelerated networking: the VF has a device, but its traffic is + // also in eth0, its master. A bond port is the same. + srv.device("enP1s1") + if err := os.Symlink("../eth0", filepath.Join(srv.root, "sys/class/net/enP1s1/master")); err != nil { + t.Fatal(err) + } + srv.net(map[string][2]uint64{"eth0": {1000, 500}, "eth1": {100, 50}, "enP1s1": {900, 400}}) + + all, err := parseFile(srv.collector(t.TempDir()), "proc/net/dev", parseNetDev) + if err != nil { + t.Fatal(err) + } + got := srv.collector(t.TempDir()).interfaces(all) + sort.Strings(got) + if fmt.Sprint(got) != "[eth0 eth1]" { + t.Errorf("interfaces() = %v, want [eth0 eth1] without the port enP1s1", got) + } +} + +func TestNewWithCountersThatCannotBeRead(t *testing.T) { + srv := newServer(t) + state := t.TempDir() + if err := os.WriteFile(filepath.Join(state, "counters.json"), []byte("{"), 0o600); err != nil { + t.Fatal(err) + } + + s, err := srv.collector(state).Sample(time.Now()) + if err != nil { + t.Fatal(err) + } + if !s.NetCountersReset { + t.Error("net_counters_reset = false, want true after counters that cannot be read") + } +} diff --git a/internal/metrics/parse.go b/internal/metrics/parse.go new file mode 100644 index 0000000..d5f6724 --- /dev/null +++ b/internal/metrics/parse.go @@ -0,0 +1,204 @@ +package metrics + +import ( + "bufio" + "bytes" + "fmt" + "strconv" + "strings" +) + +// cpuTimes are the CPU counters of /proc/stat, in clock ticks. +type cpuTimes struct { + Idle uint64 `json:"idle"` + Total uint64 `json:"total"` +} + +// parseCPU reads the first line of /proc/stat: +// +// cpu user nice system idle iowait irq softirq steal guest guest_nice +// +// Idle is idle + iowait. Total leaves out guest and guest_nice, because user +// and nice already hold them. +func parseCPU(data []byte) (cpuTimes, error) { + line, _, _ := bytes.Cut(data, []byte("\n")) + fields := strings.Fields(string(line)) + if len(fields) < 5 || fields[0] != "cpu" { + return cpuTimes{}, fmt.Errorf("/proc/stat: unexpected first line %q", line) + } + + var v [8]uint64 + for i := 0; i < len(v) && i+1 < len(fields); i++ { + n, err := strconv.ParseUint(fields[i+1], 10, 64) + if err != nil { + return cpuTimes{}, fmt.Errorf("/proc/stat: %w", err) + } + v[i] = n + } + + var t cpuTimes + for _, n := range v { + t.Total += n + } + t.Idle = v[3] + v[4] + + return t, nil +} + +// cpuPercent is the busy share of the time between two readings, 0 to 100. +func cpuPercent(prev, cur cpuTimes) float64 { + if cur.Total <= prev.Total || cur.Idle < prev.Idle { + return 0 + } + + total := float64(cur.Total - prev.Total) + idle := float64(cur.Idle - prev.Idle) + return min(max((total-idle)/total*100, 0), 100) +} + +// parseLoad reads the 1 minute load average from /proc/loadavg. +func parseLoad(data []byte) (float64, error) { + fields := strings.Fields(string(data)) + if len(fields) == 0 { + return 0, fmt.Errorf("/proc/loadavg is empty") + } + + return strconv.ParseFloat(fields[0], 64) +} + +// memory holds the values of /proc/meminfo, in bytes. +type memory struct { + total, available, swapTotal, swapFree uint64 +} + +func parseMeminfo(data []byte) (memory, error) { + values := map[string]uint64{} + s := bufio.NewScanner(bytes.NewReader(data)) + for s.Scan() { + key, rest, ok := strings.Cut(s.Text(), ":") + if !ok { + continue + } + fields := strings.Fields(rest) + if len(fields) == 0 { + continue + } + n, err := strconv.ParseUint(fields[0], 10, 64) + if err != nil { + continue + } + // The kernel shows these values in kB (KiB). + values[key] = n * 1024 + } + + for _, key := range []string{"MemTotal", "MemAvailable", "SwapTotal", "SwapFree"} { + if _, ok := values[key]; !ok { + return memory{}, fmt.Errorf("/proc/meminfo has no %s", key) + } + } + + return memory{ + total: values["MemTotal"], + available: values["MemAvailable"], + swapTotal: values["SwapTotal"], + swapFree: values["SwapFree"], + }, nil +} + +// netCounters are the received and sent bytes of one interface. +type netCounters struct { + In uint64 `json:"in"` + Out uint64 `json:"out"` +} + +// parseNetDev reads /proc/net/dev. Each interface line is +// +// name: rx_bytes rx_packets ... (8 receive fields) tx_bytes ... +func parseNetDev(data []byte) (map[string]netCounters, error) { + out := map[string]netCounters{} + s := bufio.NewScanner(bytes.NewReader(data)) + for s.Scan() { + name, rest, ok := strings.Cut(s.Text(), ":") + if !ok { + continue + } + fields := strings.Fields(rest) + if len(fields) < 9 { + continue + } + in, err1 := strconv.ParseUint(fields[0], 10, 64) + sent, err2 := strconv.ParseUint(fields[8], 10, 64) + if err1 != nil || err2 != nil { + return nil, fmt.Errorf("/proc/net/dev: bad counters for %s", strings.TrimSpace(name)) + } + out[strings.TrimSpace(name)] = netCounters{In: in, Out: sent} + } + + return out, nil +} + +// parseDefaultRoute returns the interface of the default route in +// /proc/net/route, or "". The default route has destination and mask 0: a +// VPN that routes 0.0.0.0/1 and 128.0.0.0/1 is not the default route. +// +// Iface Destination Gateway Flags RefCnt Use Metric Mask MTU Window IRTT +func parseDefaultRoute(data []byte) string { + s := bufio.NewScanner(bytes.NewReader(data)) + for s.Scan() { + fields := strings.Fields(s.Text()) + if len(fields) >= 8 && fields[1] == "00000000" && fields[7] == "00000000" { + return fields[0] + } + } + + return "" +} + +// parseOSRelease returns PRETTY_NAME from /etc/os-release. +func parseOSRelease(data []byte) string { + s := bufio.NewScanner(bytes.NewReader(data)) + for s.Scan() { + if v, ok := strings.CutPrefix(s.Text(), "PRETTY_NAME="); ok { + if unquoted, err := strconv.Unquote(v); err == nil { + return unquoted + } + return strings.Trim(v, `"'`) + } + } + + return "" +} + +// parseUptime reads the seconds since boot from /proc/uptime. +func parseUptime(data []byte) (uint64, error) { + fields := strings.Fields(string(data)) + if len(fields) == 0 { + return 0, fmt.Errorf("/proc/uptime is empty") + } + f, err := strconv.ParseFloat(fields[0], 64) + if err != nil || f < 0 { + return 0, fmt.Errorf("/proc/uptime: bad value %q", fields[0]) + } + + return uint64(f), nil +} + +// parseAptCheck reads the output of apt-check: "total;security". apt-check +// can write warnings before the result (on Ubuntu 24.04, for example for a +// source that is configured two times), so only the last line counts. +func parseAptCheck(out []byte) (total, security uint64, err error) { + lines := strings.Split(strings.TrimSpace(string(out)), "\n") + last := strings.TrimSpace(lines[len(lines)-1]) + a, b, ok := strings.Cut(last, ";") + if !ok { + return 0, 0, fmt.Errorf("apt-check: unexpected output %q", last) + } + if total, err = strconv.ParseUint(a, 10, 64); err != nil { + return 0, 0, fmt.Errorf("apt-check: %w", err) + } + if security, err = strconv.ParseUint(b, 10, 64); err != nil { + return 0, 0, fmt.Errorf("apt-check: %w", err) + } + + return total, security, nil +} diff --git a/internal/metrics/parse_test.go b/internal/metrics/parse_test.go new file mode 100644 index 0000000..ede24c7 --- /dev/null +++ b/internal/metrics/parse_test.go @@ -0,0 +1,118 @@ +package metrics + +import ( + "testing" +) + +func TestParseCPU(t *testing.T) { + // user nice system idle iowait irq softirq steal guest guest_nice + got, err := parseCPU([]byte("cpu 100 10 50 800 40 0 0 0 30 0\ncpu0 1 1 1 1 1 1 1 1 1 1\n")) + if err != nil { + t.Fatal(err) + } + // Guest time is already in user time, so it is not added again. + if got.Total != 1000 || got.Idle != 840 { + t.Errorf("parseCPU() = %+v, want total 1000 and idle 840", got) + } + + for _, bad := range []string{"", "intr 1 2 3", "cpu a b c d"} { + if _, err := parseCPU([]byte(bad)); err == nil { + t.Errorf("parseCPU(%q) = nil error, want an error", bad) + } + } +} + +func TestCPUPercent(t *testing.T) { + tests := []struct { + prev, cur cpuTimes + want float64 + }{ + {cpuTimes{Idle: 800, Total: 1000}, cpuTimes{Idle: 950, Total: 1600}, 75}, + {cpuTimes{Idle: 800, Total: 1000}, cpuTimes{Idle: 1400, Total: 1600}, 0}, + {cpuTimes{Idle: 800, Total: 1000}, cpuTimes{Idle: 800, Total: 1600}, 100}, + // No time passed, or the counters went back. + {cpuTimes{Idle: 800, Total: 1000}, cpuTimes{Idle: 800, Total: 1000}, 0}, + {cpuTimes{Idle: 800, Total: 1000}, cpuTimes{Idle: 10, Total: 20}, 0}, + } + + for _, tt := range tests { + if got := cpuPercent(tt.prev, tt.cur); got != tt.want { + t.Errorf("cpuPercent(%+v, %+v) = %v, want %v", tt.prev, tt.cur, got, tt.want) + } + } +} + +func TestParseMeminfoNeedsMemAvailable(t *testing.T) { + if _, err := parseMeminfo([]byte("MemTotal: 100 kB\nMemFree: 50 kB\nSwapTotal: 0 kB\nSwapFree: 0 kB\n")); err == nil { + t.Error("parseMeminfo() = nil error, want an error without MemAvailable") + } +} + +func TestParseNetDev(t *testing.T) { + data := "Inter-| Receive | Transmit\n face |bytes packets|bytes\n" + + " eth0: 1234 5 0 0 0 0 0 0 5678 6 0 0 0 0 0 0\n" + + " lo:1 1 0 0 0 0 0 0 2 2 0 0 0 0 0 0\n" + + got, err := parseNetDev([]byte(data)) + if err != nil { + t.Fatal(err) + } + if got["eth0"] != (netCounters{In: 1234, Out: 5678}) || got["lo"] != (netCounters{In: 1, Out: 2}) || len(got) != 2 { + t.Errorf("parseNetDev() = %+v", got) + } +} + +func TestParseDefaultRoute(t *testing.T) { + const header = "Iface\tDestination\tGateway \tFlags\tRefCnt\tUse\tMetric\tMask\t\tMTU\tWindow\tIRTT\n" + data := header + + "eth1\t0A000000\t00000000\t0001\t0\t0\t0\t000000FF\t0\t0\t0\n" + + "eth0\t00000000\t01C0A8C0\t0003\t0\t0\t100\t00000000\t0\t0\t0\n" + if got := parseDefaultRoute([]byte(data)); got != "eth0" { + t.Errorf("parseDefaultRoute() = %q, want eth0", got) + } + + // OpenVPN def1: 0.0.0.0/1 on tun0 is not the default route. + vpn := header + "tun0\t00000000\t0100080A\t0003\t0\t0\t0\t00000080\t0\t0\t0\n" + if got := parseDefaultRoute([]byte(vpn)); got != "" { + t.Errorf("parseDefaultRoute() = %q, want no default route for 0.0.0.0/1", got) + } + if got := parseDefaultRoute([]byte("Iface\tDestination\n")); got != "" { + t.Errorf("parseDefaultRoute() = %q, want no interface", got) + } +} + +func TestParseOSRelease(t *testing.T) { + tests := map[string]string{ + "NAME=\"Ubuntu\"\nPRETTY_NAME=\"Ubuntu 24.04.1 LTS\"\n": "Ubuntu 24.04.1 LTS", + "PRETTY_NAME='Debian GNU/Linux 12'\n": "Debian GNU/Linux 12", + "PRETTY_NAME=Alpine\n": "Alpine", + "NAME=x\n": "", + } + for in, want := range tests { + if got := parseOSRelease([]byte(in)); got != want { + t.Errorf("parseOSRelease(%q) = %q, want %q", in, got, want) + } + } +} + +func TestParseUptime(t *testing.T) { + if got, err := parseUptime([]byte("350735.47 234388.90\n")); err != nil || got != 350735 { + t.Errorf("parseUptime() = %d, %v; want 350735", got, err) + } + if _, err := parseUptime([]byte("")); err == nil { + t.Error("parseUptime(\"\") = nil error, want an error") + } +} + +func TestParseAptCheck(t *testing.T) { + for _, out := range []string{"33;6", "33;6\n", "Warning: W:something; else\n33;6"} { + if total, security, err := parseAptCheck([]byte(out)); err != nil || total != 33 || security != 6 { + t.Errorf("parseAptCheck(%q) = %d, %d, %v; want 33, 6", out, total, security, err) + } + } + for _, bad := range []string{"", "33", "a;b", "1;x"} { + if _, _, err := parseAptCheck([]byte(bad)); err == nil { + t.Errorf("parseAptCheck(%q) = nil error, want an error", bad) + } + } +} diff --git a/internal/metrics/sys_linux.go b/internal/metrics/sys_linux.go new file mode 100644 index 0000000..2686e0d --- /dev/null +++ b/internal/metrics/sys_linux.go @@ -0,0 +1,34 @@ +package metrics + +import ( + "golang.org/x/sys/unix" +) + +// statfs returns the size and the used space of the file system that holds +// path. Used space is (blocks − free blocks): the space that is reserved for +// root is used space, because the sites cannot use it. +func statfs(path string) (total, used uint64, err error) { + var st unix.Statfs_t + if err := unix.Statfs(path, &st); err != nil { + return 0, 0, err + } + + // The block counts are in units of the fragment size (f_frsize), the + // same as df uses. + size := uint64(st.Frsize) + if size == 0 { + size = uint64(st.Bsize) + } + + return st.Blocks * size, (st.Blocks - st.Bfree) * size, nil +} + +// kernelRelease returns the kernel release, as "uname -r" shows it. +func kernelRelease() string { + var u unix.Utsname + if err := unix.Uname(&u); err != nil { + return "" + } + + return unix.ByteSliceToString(u.Release[:]) +} diff --git a/internal/metrics/sys_other.go b/internal/metrics/sys_other.go new file mode 100644 index 0000000..af62166 --- /dev/null +++ b/internal/metrics/sys_other.go @@ -0,0 +1,19 @@ +//go:build !linux + +package metrics + +import ( + "errors" + "runtime" +) + +// The agent measures only Linux servers. On other systems, for example a +// developer's Mac, the disk and the kernel are not measured. + +func statfs(string) (total, used uint64, err error) { + return 0, 0, errors.New("disk measurement is not supported on " + runtime.GOOS) +} + +func kernelRelease() string { + return "" +}