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
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
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, 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`).

Expand Down
15 changes: 15 additions & 0 deletions internal/agent/clean.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ const (
maxVersionLen = 32
maxStatusTextLen = 255
maxArchLen = 16
maxCPUCount = 4096
maxEventNameLen = 64
maxErrorLen = 2000
)
Expand Down Expand Up @@ -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
}

Expand Down
29 changes: 29 additions & 0 deletions internal/agent/outbox_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
14 changes: 14 additions & 0 deletions internal/agent/wire/wire.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Expand Down
98 changes: 98 additions & 0 deletions internal/dockerapi/dockerapi.go
Original file line number Diff line number Diff line change
@@ -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
}
57 changes: 57 additions & 0 deletions internal/dockerapi/dockerapi_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
45 changes: 45 additions & 0 deletions internal/metrics/docker.go
Original file line number Diff line number Diff line change
@@ -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
}
}
Loading
Loading