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: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

Easy CLI tool for servers managed by FlyWP.

Conforms to the FlyWP monitoring agent contract v0.3.1.
Conforms to the FlyWP monitoring agent contract v0.4.0.

## Installation

Expand Down Expand Up @@ -103,7 +103,7 @@ All arguments after the WP-CLI command (or after the command for `fly exec`) go

### Monitoring agent

`fly agent run` is the FlyWP monitoring agent. It runs all the time under systemd (`fly-agent.service`, as the server user, not root), and FlyWP installs it. Each minute it measures CPU, load, memory, swap, disk and network traffic, and it sends the values and the server status (restart needed, waiting updates, OS, kernel, uptime) to FlyWP. It keeps unsent data on disk for up to 24 hours. FlyWP can update and restart the agent through it, without SSH. The agent does not need Docker.
`fly agent run` is the FlyWP monitoring agent. It runs all the time under systemd (`fly-agent.service`, as the server user, not root), and FlyWP installs it. Each minute it measures CPU, load, memory, swap, disk and network traffic, the pressure (PSI) and the disk activity. It reads the server each 10 seconds, so each minute also has its peaks. It sends the values and the server status (restart needed, waiting updates, OS, kernel, uptime) to FlyWP. It keeps unsent data on disk for up to 24 hours. FlyWP can update and restart the agent through it, without SSH. The agent does not need Docker.

It reads `FLY_AGENT_URL` (https), `FLY_AGENT_TOKEN` and `FLY_AGENT_SERVER_ID` from `/etc/fly/agent.env`, and keeps its state in `STATE_DIRECTORY` (`/var/lib/fly-agent`).

Expand Down
2 changes: 2 additions & 0 deletions internal/agent/clean.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,8 @@ func cleanSample(s wire.Sample) wire.Sample {
}
for _, v := range []**uint64{
&s.MemoryUsedMaxBytes, &s.SwapUsedMaxBytes, &s.NetInMaxBytesPerSecond, &s.NetOutMaxBytesPerSecond,
&s.DiskReadBytes, &s.DiskWriteBytes, &s.DiskReadOps, &s.DiskWriteOps,
&s.DiskReadMaxBytesPerSecond, &s.DiskWriteMaxBytesPerSecond, &s.DiskReadMaxOpsPerSecond, &s.DiskWriteMaxOpsPerSecond,
} {
*v = clampIntPtr(*v)
}
Expand Down
2 changes: 1 addition & 1 deletion internal/agent/config.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
// Package agent is the FlyWP monitoring agent: the long-running mode of fly
// that "fly agent run" starts. It follows the FlyWP monitoring agent
// contract v0.3.1.
// contract v0.4.0.
package agent

import (
Expand Down
13 changes: 12 additions & 1 deletion internal/agent/wire/wire.go
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// Package wire holds the JSON bodies of the FlyWP monitoring agent contract
// v0.2.1: the requests that the agent sends and the replies that it reads.
// v0.4.0: the requests that the agent sends and the replies that it reads.
package wire

import (
Expand Down Expand Up @@ -63,6 +63,17 @@ type Sample struct {
MemoryPressureMaxPercent *float64 `json:"memory_pressure_max_percent"`
IOPressurePercent *float64 `json:"io_pressure_percent"`
IOPressureMaxPercent *float64 `json:"io_pressure_max_percent"`

// The disk activity of the minute, and its peaks each second within the
// minute (contract v0.4.0). nil (JSON null) means "not known".
DiskReadBytes *uint64 `json:"disk_read_bytes"`
DiskWriteBytes *uint64 `json:"disk_write_bytes"`
DiskReadOps *uint64 `json:"disk_read_ops"`
DiskWriteOps *uint64 `json:"disk_write_ops"`
DiskReadMaxBytesPerSecond *uint64 `json:"disk_read_max_bytes_per_second"`
DiskWriteMaxBytesPerSecond *uint64 `json:"disk_write_max_bytes_per_second"`
DiskReadMaxOpsPerSecond *uint64 `json:"disk_read_max_ops_per_second"`
DiskWriteMaxOpsPerSecond *uint64 `json:"disk_write_max_ops_per_second"`
}

// MetricsReply is the reply to POST /agent/v1/metrics.
Expand Down
99 changes: 99 additions & 0 deletions internal/metrics/disk.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
package metrics

import (
"os"
"slices"

"github.com/flywp/server-cli/internal/agent/wire"
)

// sectorSize is the unit of the sector counts in /proc/diskstats, for each
// disk, whatever the sector size of the disk.
const sectorSize = 512

// readDisks returns the counters of the disks that have a hardware device: a
// name in /sys/block with a device link, for example vda, sda or nvme0n1.
// Partitions are not in /sys/block, and loop, ram, zram, device mapper and
// software RAID disks have no device: their activity is already in a hardware
// disk. It returns nil when /proc/diskstats cannot be read.
func (c *Collector) readDisks() map[string]diskCounters {
all, err := parseFile(c, "proc/diskstats", parseDiskstats)
if err != nil {
c.log.Debug("reading the disk activity", "error", err)
return nil
}

disks := map[string]diskCounters{}
for name, d := range all {
if _, err := os.Lstat(c.file("sys/block", name, "device")); err == nil {
disks[name] = d
}
}
return disks
}

// diskDelta returns the activity between two readings. It adds the disks
// that both readings have, so a new or a removed disk makes no spike. ok is
// false when the activity is not known: a reading without disks, a reboot, a
// reading older than maxAge, or a counter that went back.
func diskDelta(a, b reading) (d diskCounters, ok bool) {
if len(a.Disks) == 0 || len(b.Disks) == 0 || a.BootID != b.BootID {
return diskCounters{}, false
}
if age := b.At.Sub(a.At); age <= 0 || age > maxAge {
return diskCounters{}, false
}

for name, cur := range b.Disks {
prev, found := a.Disks[name]
if !found {
continue
}
if cur.ReadOps < prev.ReadOps || cur.ReadSectors < prev.ReadSectors ||
cur.WriteOps < prev.WriteOps || cur.WriteSectors < prev.WriteSectors {
return diskCounters{}, false
}
d.ReadOps += cur.ReadOps - prev.ReadOps
d.ReadSectors += cur.ReadSectors - prev.ReadSectors
d.WriteOps += cur.WriteOps - prev.WriteOps
d.WriteSectors += cur.WriteSectors - prev.WriteSectors
}

return d, true
}

// setDiskActivity sets the disk activity of the minute from prev to cur in s,
// and the peaks of the windows in all. The eight fields stay nil when the
// activity of the minute is not known, or when a window has a counter that
// went back: then the delta of the minute is not the real activity either.
func setDiskActivity(s *wire.Sample, prev *reading, all []reading, cur reading) {
if prev == nil {
return
}
d, ok := diskDelta(*prev, cur)
if !ok {
return
}
readBytes, writeBytes := d.ReadSectors*sectorSize, d.WriteSectors*sectorSize

// The control plane reads the value of the minute as the value / 60.
peak := [4]uint64{readBytes / 60, writeBytes / 60, d.ReadOps / 60, d.WriteOps / 60}
// A reading without disks is left out: its windows join.
all = slices.DeleteFunc(slices.Clone(all), func(r reading) bool { return len(r.Disks) == 0 })
for i := 1; i < len(all); i++ {
a, b := all[i-1], all[i]
w, ok := diskDelta(a, b)
if !ok {
return
}
secs := b.At.Sub(a.At).Seconds()
for j, v := range []uint64{w.ReadSectors * sectorSize, w.WriteSectors * sectorSize, w.ReadOps, w.WriteOps} {
peak[j] = max(peak[j], uint64(float64(v)/secs))
}
}

s.DiskReadBytes, s.DiskWriteBytes = &readBytes, &writeBytes
s.DiskReadOps, s.DiskWriteOps = &d.ReadOps, &d.WriteOps
s.DiskReadMaxBytesPerSecond, s.DiskWriteMaxBytesPerSecond = &peak[0], &peak[1]
s.DiskReadMaxOpsPerSecond, s.DiskWriteMaxOpsPerSecond = &peak[2], &peak[3]
}
Loading
Loading