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
4 changes: 3 additions & 1 deletion cmd/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"syscall"

"github.com/flywp/server-cli/internal/agent"
"github.com/flywp/server-cli/internal/metrics"
"github.com/spf13/cobra"
)

Expand All @@ -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))
},
}

Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
)

Expand All @@ -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
)
322 changes: 322 additions & 0 deletions internal/metrics/metrics.go
Original file line number Diff line number Diff line change
@@ -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)
}
34 changes: 34 additions & 0 deletions internal/metrics/metrics_linux_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
Loading
Loading