diff --git a/README.md b/README.md index bfd7f11..5c2fdaa 100644 --- a/README.md +++ b/README.md @@ -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, 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. +`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, CPU count, Docker state and version) 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. When Docker runs, the agent reads its socket with two requests only: `GET /version` and `GET /containers/json`. 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`). diff --git a/internal/agent/clean.go b/internal/agent/clean.go index 782c39b..46f5bd6 100644 --- a/internal/agent/clean.go +++ b/internal/agent/clean.go @@ -18,6 +18,7 @@ const ( maxVersionLen = 32 maxStatusTextLen = 255 maxArchLen = 16 + maxCPUCount = 4096 maxEventNameLen = 64 maxErrorLen = 2000 ) @@ -55,6 +56,20 @@ func cleanStatus(s wire.Status) wire.Status { s.UpdatesTotal = clampIntPtr(s.UpdatesTotal) s.UpdatesSecurity = clampIntPtr(s.UpdatesSecurity) s.UptimeSeconds = clampInt(s.UptimeSeconds) + if s.CPUCount != nil && (*s.CPUCount < 1 || *s.CPUCount > maxCPUCount) { + s.CPUCount = nil + } + if s.DockerStatus != nil { + switch *s.DockerStatus { + case wire.DockerRunning, wire.DockerNotRunning, wire.DockerNotInstalled: + default: + s.DockerStatus = nil + } + } + if s.DockerVersion != nil { + v := truncate(*s.DockerVersion, maxVersionLen) + s.DockerVersion = &v + } return s } diff --git a/internal/agent/outbox_test.go b/internal/agent/outbox_test.go index 15da7e9..8fdcecd 100644 --- a/internal/agent/outbox_test.go +++ b/internal/agent/outbox_test.go @@ -213,6 +213,35 @@ func TestQueuedSampleOfAnOlderAgentSendsNullPeaks(t *testing.T) { } } +func TestCleanDockerStatus(t *testing.T) { + zero, many, four := 0, 5000, 4 + status, version := "paused", strings.Repeat("9", 40) + s := cleanStatus(wire.Status{CPUCount: &zero, DockerStatus: &status, DockerVersion: &version}) + if s.CPUCount != nil || s.DockerStatus != nil || s.DockerVersion == nil || len(*s.DockerVersion) != maxVersionLen { + t.Errorf("cleanStatus() = %v, %v, %v; want null, null and %d characters", s.CPUCount, s.DockerStatus, s.DockerVersion, maxVersionLen) + } + if s := cleanStatus(wire.Status{CPUCount: &many}); s.CPUCount != nil { + t.Errorf("cpu_count = %d, want null above %d", *s.CPUCount, maxCPUCount) + } + + running := wire.DockerRunning + s = cleanStatus(wire.Status{CPUCount: &four, DockerStatus: &running}) + if s.CPUCount == nil || *s.CPUCount != 4 || s.DockerStatus == nil || *s.DockerStatus != wire.DockerRunning { + t.Errorf("cleanStatus() = %v, %v; want 4 and running", s.CPUCount, s.DockerStatus) + } + + // Not known is null on the wire. + data, err := json.Marshal(cleanStatus(wire.Status{})) + if err != nil { + t.Fatal(err) + } + for _, field := range []string{"cpu_count", "docker_status", "docker_version"} { + if !strings.Contains(string(data), `"`+field+`":null`) { + t.Errorf("status 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/wire/wire.go b/internal/agent/wire/wire.go index 0f8839d..0247254 100644 --- a/internal/agent/wire/wire.go +++ b/internal/agent/wire/wire.go @@ -27,8 +27,22 @@ type Status struct { Kernel string `json:"kernel"` UptimeSeconds uint64 `json:"uptime_seconds"` Arch string `json:"arch"` + + // CPUCount is the number of CPUs that are online. DockerStatus is one of + // the Docker* values, and DockerVersion is set only when Docker runs + // (contract v0.5.0). nil (JSON null) means "not known". + CPUCount *int `json:"cpu_count"` + DockerStatus *string `json:"docker_status"` + DockerVersion *string `json:"docker_version"` } +// The values of Status.DockerStatus. +const ( + DockerRunning = "running" + DockerNotRunning = "not_running" + DockerNotInstalled = "not_installed" +) + // Sample holds the measurements of one minute. type Sample struct { RecordedAt time.Time `json:"recorded_at"` diff --git a/internal/dockerapi/dockerapi.go b/internal/dockerapi/dockerapi.go new file mode 100644 index 0000000..97a8ba0 --- /dev/null +++ b/internal/dockerapi/dockerapi.go @@ -0,0 +1,98 @@ +// Package dockerapi reads the Docker Engine API over its unix socket, with the +// Go standard library only. It sends two requests, GET /version and +// GET /containers/json, and never another one: access to the socket is root +// access on the server (FlyWP monitoring agent contract v0.5.0). +package dockerapi + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net" + "net/http" + "time" +) + +// DefaultTimeout limits each request (contract v0.5.0). +const DefaultTimeout = 5 * time.Second + +// maxBody limits the size of an answer. +const maxBody = 16 << 20 + +// Client reads the Docker Engine API. Use New. +type Client struct { + // Timeout limits each request. Tests make it shorter. + Timeout time.Duration + + hc *http.Client +} + +// New returns a client for the socket at path, for example +// /var/run/docker.sock. +func New(path string) *Client { + var d net.Dialer + return &Client{ + Timeout: DefaultTimeout, + hc: &http.Client{ + Transport: &http.Transport{ + DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) { + return d.DialContext(ctx, "unix", path) + }, + MaxIdleConns: 1, + }, + CheckRedirect: func(*http.Request, []*http.Request) error { + return http.ErrUseLastResponse + }, + }, + } +} + +// Container is a running container. +type Container struct { + ID string `json:"Id"` + Labels map[string]string `json:"Labels"` +} + +// Version returns the version of the Docker Engine, for example 29.7.1. +func (c *Client) Version(ctx context.Context) (string, error) { + var v struct { + Version string `json:"Version"` + } + if err := c.get(ctx, "/version", &v); err != nil { + return "", err + } + return v.Version, nil +} + +// Containers returns the running containers. +func (c *Client) Containers(ctx context.Context) ([]Container, error) { + var cs []Container + if err := c.get(ctx, "/containers/json", &cs); err != nil { + return nil, err + } + return cs, nil +} + +func (c *Client) get(ctx context.Context, path string, out any) error { + ctx, cancel := context.WithTimeout(ctx, c.Timeout) + defer cancel() + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, "http://docker"+path, nil) + if err != nil { + return err + } + resp, err := c.hc.Do(req) + if err != nil { + return err + } + defer func() { _ = resp.Body.Close() }() + + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("docker %s: %s", path, resp.Status) + } + if err := json.NewDecoder(io.LimitReader(resp.Body, maxBody)).Decode(out); err != nil { + return fmt.Errorf("docker %s: %w", path, err) + } + return nil +} diff --git a/internal/dockerapi/dockerapi_test.go b/internal/dockerapi/dockerapi_test.go new file mode 100644 index 0000000..a500470 --- /dev/null +++ b/internal/dockerapi/dockerapi_test.go @@ -0,0 +1,57 @@ +package dockerapi + +import ( + "context" + "errors" + "net/http" + "testing" + "time" + + "github.com/flywp/server-cli/internal/testutil" +) + +var engineHandler = testutil.EngineHandler("29.7.1", + `[{"Id":"abc","Names":["/x"],"Labels":{"com.docker.compose.project.working_dir":"/home/fly/example.com"}},{"Id":"def","Labels":{}}]`) + +func TestVersionAndContainers(t *testing.T) { + e := testutil.NewEngine(t, engineHandler) + c := New(e.Socket) + + v, err := c.Version(context.Background()) + if err != nil || v != "29.7.1" { + t.Errorf("Version() = %q, %v; want 29.7.1", v, err) + } + cs, err := c.Containers(context.Background()) + if err != nil || len(cs) != 2 || cs[0].ID != "abc" || cs[0].Labels["com.docker.compose.project.working_dir"] != "/home/fly/example.com" { + t.Errorf("Containers() = %+v, %v", cs, err) + } + if got := e.Requests(); len(got) != 2 || got[0] != "GET /version" || got[1] != "GET /containers/json" { + t.Errorf("requests = %v, want only GET /version and GET /containers/json", got) + } +} + +func TestErrorAnswer(t *testing.T) { + e := testutil.NewEngine(t, func(w http.ResponseWriter, _ *http.Request) { + http.Error(w, "boom", http.StatusInternalServerError) + }) + if _, err := New(e.Socket).Version(context.Background()); err == nil { + t.Error("Version() = nil error, want the 500") + } +} + +func TestTimeout(t *testing.T) { + release := make(chan struct{}) + e := testutil.NewEngine(t, func(http.ResponseWriter, *http.Request) { <-release }) + defer close(release) + + c := New(e.Socket) + c.Timeout = 50 * time.Millisecond + start := time.Now() + _, err := c.Version(context.Background()) + if !errors.Is(err, context.DeadlineExceeded) { + t.Errorf("Version() = %v, want a timeout", err) + } + if d := time.Since(start); d > 2*time.Second { + t.Errorf("Version() took %v, want the timeout", d) + } +} diff --git a/internal/metrics/docker.go b/internal/metrics/docker.go new file mode 100644 index 0000000..3260108 --- /dev/null +++ b/internal/metrics/docker.go @@ -0,0 +1,45 @@ +package metrics + +import ( + "context" + "errors" + "io/fs" + "os" + "syscall" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +// dockerStatus returns the state and the version of Docker, from GET /version +// on the Docker socket (contract v0.5.0): +// +// - an answer: running +// - no socket, or the connection is refused: not_running when dockerd +// exists, else not_installed +// - any other result, for example "permission denied" or no answer in 5 +// seconds: nil, because the agent cannot tell +func (c *Collector) dockerStatus(ctx context.Context) (status, version *string) { + v, err := c.docker.Version(ctx) + switch { + case err == nil: + c.dockerWarned = false + s := wire.DockerRunning + if v == "" { + return &s, nil + } + return &s, &v + case errors.Is(err, fs.ErrNotExist) || errors.Is(err, syscall.ECONNREFUSED): + c.dockerWarned = false + s := wire.DockerNotInstalled + if _, err := os.Stat(c.file("usr/bin/dockerd")); err == nil { + s = wire.DockerNotRunning + } + return &s, nil + default: + if !c.dockerWarned { + c.dockerWarned = true + c.log.Warn("cannot tell whether Docker runs; sending null", "error", err) + } + return nil, nil + } +} diff --git a/internal/metrics/docker_test.go b/internal/metrics/docker_test.go new file mode 100644 index 0000000..8cc6bf0 --- /dev/null +++ b/internal/metrics/docker_test.go @@ -0,0 +1,131 @@ +package metrics + +import ( + "context" + "net" + "net/http" + "os" + "path/filepath" + "testing" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" + "github.com/flywp/server-cli/internal/dockerapi" + "github.com/flywp/server-cli/internal/testutil" +) + +// shortDir returns a folder with a short path, for a unix socket: its path +// has at most 104 bytes on macOS. +func shortDir(t *testing.T) string { + t.Helper() + dir, err := os.MkdirTemp("", "dk") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.RemoveAll(dir) }) + return dir +} + +func TestCPUCount(t *testing.T) { + srv := newServer(t) + srv.write("proc/stat", "cpu 1 2 3 4 5 6 7 8 0 0\ncpu0 1 2 3 4 5 6 7 8 0 0\ncpu1 1 2 3 4 5 6 7 8 0 0\ncpu2 1 2 3 4 5 6 7 8 0 0\ncpu3 1 2 3 4 5 6 7 8 0 0\nintr 1 2 3\nctxt 5\ncpufreq 1\n") + s := srv.collector(t.TempDir()).Status(context.Background()) + if s.CPUCount == nil || *s.CPUCount != 4 { + t.Errorf("cpu_count = %v, want 4", ptr(s.CPUCount)) + } +} + +func TestDockerStatus(t *testing.T) { + engine := testutil.NewEngine(t, testutil.EngineHandler("29.7.1", "[]")) + + refused := filepath.Join(shortDir(t), "docker.sock") + l, err := net.Listen("unix", refused) + if err != nil { + t.Fatal(err) + } + l.(*net.UnixListener).SetUnlinkOnClose(false) + if err := l.Close(); err != nil { + t.Fatal(err) + } + + tests := []struct { + name string + socket string + dockerd bool + wantStatus any + wantVersion any + }{ + {"running", engine.Socket, true, wire.DockerRunning, "29.7.1"}, + {"no socket, with dockerd", "/nonexistent/docker.sock", true, wire.DockerNotRunning, "null"}, + {"no socket, without dockerd", "/nonexistent/docker.sock", false, wire.DockerNotInstalled, "null"}, + {"connection refused", refused, true, wire.DockerNotRunning, "null"}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + srv := newServer(t) + if tt.dockerd { + srv.write("usr/bin/dockerd", "") + } + c := srv.collector(t.TempDir()) + c.docker = dockerapi.New(tt.socket) + + s := c.Status(context.Background()) + if ptr(s.DockerStatus) != tt.wantStatus || ptr(s.DockerVersion) != tt.wantVersion { + t.Errorf("docker = %v %v, want %v %v", ptr(s.DockerStatus), ptr(s.DockerVersion), tt.wantStatus, tt.wantVersion) + } + }) + } +} + +func TestDockerStatusIsNotKnown(t *testing.T) { + t.Run("permission denied", func(t *testing.T) { + if os.Getuid() == 0 { + t.Skip("root can use any socket") + } + engine := testutil.NewEngine(t, testutil.EngineHandler("29.7.1", "[]")) + if err := os.Chmod(engine.Socket, 0); err != nil { + t.Fatal(err) + } + srv := newServer(t) + srv.write("usr/bin/dockerd", "") + c := srv.collector(t.TempDir()) + c.docker = dockerapi.New(engine.Socket) + + if s := c.Status(context.Background()); s.DockerStatus != nil || s.DockerVersion != nil { + t.Errorf("docker = %v %v, want null: the agent cannot tell", ptr(s.DockerStatus), ptr(s.DockerVersion)) + } + }) + + t.Run("no answer", func(t *testing.T) { + release := make(chan struct{}) + engine := testutil.NewEngine(t, func(http.ResponseWriter, *http.Request) { <-release }) + defer close(release) + + srv := newServer(t) + c := srv.collector(t.TempDir()) + c.docker = dockerapi.New(engine.Socket) + c.docker.Timeout = 50 * time.Millisecond + + s := c.Status(context.Background()) + if s.DockerStatus != nil || s.DockerVersion != nil { + t.Errorf("docker = %v %v, want null", ptr(s.DockerStatus), ptr(s.DockerVersion)) + } + // The other values are still there. + if s.OS == "" || s.CPUCount == nil { + t.Errorf("status = %+v, want the other values", s) + } + }) +} + +func TestDockerStatusSendsOnlyVersion(t *testing.T) { + engine := testutil.NewEngine(t, testutil.EngineHandler("29.7.1", "[]")) + srv := newServer(t) + c := srv.collector(t.TempDir()) + c.docker = dockerapi.New(engine.Socket) + + c.Status(context.Background()) + if got := engine.Requests(); len(got) != 1 || got[0] != "GET /version" { + t.Errorf("requests = %v, want only GET /version", got) + } +} diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index 8b32869..c8c1313 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -20,6 +20,7 @@ import ( "github.com/flywp/server-cli/internal/agent" "github.com/flywp/server-cli/internal/agent/wire" + "github.com/flywp/server-cli/internal/dockerapi" "github.com/flywp/server-cli/internal/statefile" ) @@ -71,6 +72,7 @@ type Collector struct { statfs func(path string) (total, used uint64, err error) release func() string aptCheck func(ctx context.Context) ([]byte, error) + docker *dockerapi.Client // 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. @@ -85,6 +87,9 @@ type Collector struct { // and noPSI after the warning that the kernel has no PSI. noInterface bool noPSI bool + // dockerWarned is true after the warning that the state of Docker is + // not known, until it is known again. + dockerWarned bool updatesAt time.Time updatesKnown bool @@ -104,6 +109,7 @@ func New(root, stateDir string, log *slog.Logger) *Collector { aptCheck: runAptCheck, minFirst: minFirstMinute, } + c.docker = dockerapi.New(c.file("var/run/docker.sock")) // Take a reading now, so that the first sample has a CPU value for the // time since the start. The saved reading comes before it only when it is @@ -319,10 +325,10 @@ func (c *Collector) interfaces(all map[string]netCounters) []string { 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 { +// Status describes the server now. A value that cannot be read stays empty, +// 0 or nil, and the problem goes to the log. The update counts come from the +// last Sample. +func (c *Collector) Status(ctx 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 { @@ -341,6 +347,14 @@ func (c *Collector) Status(context.Context) wire.Status { c.log.Warn("reading the uptime", "error", err) } + if n, err := parseFile(c, "proc/stat", parseCPUCount); err == nil { + s.CPUCount = &n + } else { + c.log.Warn("counting the CPUs", "error", err) + } + + s.DockerStatus, s.DockerVersion = c.dockerStatus(ctx) + // Without any count, the counts are not known: null, not a false 0. if c.updatesKnown { total, security := c.updatesTotal, c.updatesSecurity diff --git a/internal/metrics/parse.go b/internal/metrics/parse.go index 34512d4..2c64d54 100644 --- a/internal/metrics/parse.go +++ b/internal/metrics/parse.go @@ -265,3 +265,22 @@ func parseDiskstats(data []byte) (map[string]diskCounters, error) { return out, nil } + +// parseCPUCount counts the CPUs that are online: the cpuN lines of /proc/stat. +func parseCPUCount(data []byte) (int, error) { + n := 0 + s := bufio.NewScanner(bytes.NewReader(data)) + for s.Scan() { + name, _, _ := strings.Cut(s.Text(), " ") + if rest, ok := strings.CutPrefix(name, "cpu"); ok && rest != "" { + if _, err := strconv.Atoi(rest); err == nil { + n++ + } + } + } + if n == 0 { + return 0, fmt.Errorf("/proc/stat has no cpuN line") + } + + return n, nil +} diff --git a/internal/testutil/fakedocker.go b/internal/testutil/fakedocker.go index ce46d5f..c7ec06f 100644 --- a/internal/testutil/fakedocker.go +++ b/internal/testutil/fakedocker.go @@ -1,5 +1,6 @@ // Package testutil provides helpers for tests that run fly against a fake -// docker command instead of a real Docker installation. +// docker command or a fake Docker Engine API instead of a real Docker +// installation. package testutil import ( diff --git a/internal/testutil/fakeengine.go b/internal/testutil/fakeengine.go new file mode 100644 index 0000000..7df39a6 --- /dev/null +++ b/internal/testutil/fakeengine.go @@ -0,0 +1,69 @@ +package testutil + +import ( + "net" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "sync" + "testing" +) + +// Engine is a fake Docker Engine API on a unix socket. It records the +// requests. +type Engine struct { + Socket string + + mu sync.Mutex + requests []string +} + +// NewEngine starts a fake engine that answers with handler. The socket path is +// short: a unix socket path has at most 104 bytes on macOS. +func NewEngine(t *testing.T, handler http.HandlerFunc) *Engine { + t.Helper() + dir, err := os.MkdirTemp("", "dk") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.RemoveAll(dir) }) + + e := &Engine{Socket: filepath.Join(dir, "docker.sock")} + l, err := net.Listen("unix", e.Socket) + if err != nil { + t.Fatal(err) + } + srv := httptest.NewUnstartedServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + e.mu.Lock() + e.requests = append(e.requests, r.Method+" "+r.URL.Path) + e.mu.Unlock() + handler(w, r) + })) + srv.Listener = l + srv.Start() + t.Cleanup(srv.Close) + return e +} + +// Requests returns the method and the path of each request. +func (e *Engine) Requests() []string { + e.mu.Lock() + defer e.mu.Unlock() + return append([]string(nil), e.requests...) +} + +// EngineHandler answers GET /version with version, and GET /containers/json +// with containers (JSON). +func EngineHandler(version, containers string) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/version": + _, _ = w.Write([]byte(`{"Platform":{"Name":"Docker Engine - Community"},"Version":"` + version + `","ApiVersion":"1.52"}`)) + case "/containers/json": + _, _ = w.Write([]byte(containers)) + default: + http.NotFound(w, r) + } + } +}