Skip to content
Approvalens

Server agent · open source

The server agent, in the open

Before you run anything as root on your server you should be able to read it. This page lists exactly what the agent reads and sends, shows its source code, and explains how to build it yourself and check that you get the same binary we ship.

  • MIT licence
  • Go standard library only, no dependencies
  • Static binary, about 6 MB
  • Reproducible builds

What it reads, and when it is sent

The agent samples locally every 60 seconds. Very little of that goes over the network: see the schedule below.

  • CPU usage, load average, memory and swap use, swap-in/out rateRead from:/proc/stat, /proc/loadavg, /proc/meminfo, /proc/vmstatSent:Heartbeat (CPU, memory) every 5 min; per-minute average and peak in the 15-minute rollup
  • Disk and inode usage per real filesystem (ext4, xfs, btrfs, zfs…; not tmpfs, overlays, snaps or network mounts), fill rate over the last 6 hours and “full in N hours”Read from:/proc/self/mountinfo, statfs()Sent:Fullest disk and soonest ETA in the heartbeat; all filesystems in the rollup
  • Network receive / transmit rate (physical interfaces; not lo, docker, veth or bridges)Read from:/proc/net/devSent:Rollup
  • Hostname, OS name, kernel, CPU count, architecture, boot time, uptimeRead from:/proc/sys/kernel, /etc/os-release, /proc/uptimeSent:Rollup
  • Top 10 processes by CPU and by memory: pid, name, user, executable path, command line cut to 200 characters with password=…, --token …, user:pass@ and mysql -p… maskedRead from:/proc/PID/stat, cmdline, exeSent:Rollup
  • Listening TCP ports and UDP services (protocol, port, loopback-only or not) with the owning processRead from:/proc/net/tcp, tcp6, udp, udp6, /proc/PID/fdSent:Full list once after install; afterwards only what was added or removed
  • Active crontab linesRead from:/etc/crontab, /etc/cron.d, /var/spool/cronSent:Only changes
  • SSH keys of root and of accounts with a login shell: type, SHA-256 fingerprint and comment (never the key)Read from:~/.ssh/authorized_keys, authorized_keys2Sent:Only changes
  • setuid / setgid files (path, mode, owner), every 6 hoursRead from:/bin, /sbin, /usr/bin, /usr/sbin, /usr/lib, /usr/libexec, /usr/local, /opt, /tmp, /var/tmp, /dev/shmSent:Only changes
  • Only if you configure web roots: names of new, changed and removed .php, .phtml, .phar, .inc, .js, .mjs, .html, .htm, .shtml, .htaccess and .user.ini files, every 15 minutes. Compared by modification time and size; the content is never readRead from:the folders you list in web_rootsSent:Only changes (at most 50 names per list)
  • Once a day: a hash of each list above, so we can notice that our copy driftedRead from:–Sent:Rollup, once a day

What goes over the network

  1. Heartbeat · every 5 minutes · about 150 bytes

    Time, agent version, uptime, CPU %, memory %, the fullest disk and its projected-full ETA, the rules that are tripped now. Fifteen minutes without one is the “server or agent silent” alert.

  2. Rollup · every 15 minutes · about 2 KB

    Per-minute averages and peaks since the previous rollup (CPU, memory, swap, load, network, fullest disk), filesystems, the top processes. This is what the charts show.

  3. Event · at once, at most one a minute

    When a local rule trips (disk at 90% or full within 48 hours, inodes at 90%, memory at 95% for 10 minutes, heavy swapping for 10 minutes, load above twice the cores for 15 minutes, a security signal) or a list changes (ports, crontab, SSH keys, setuid, web roots), with the evidence.

  4. Inventory · rarely

    The full lists of ports, crontab lines, SSH key fingerprints and setuid files: once after install, and again only if our daily hash check says our copy differs.

Measured: about 230 KB of data a day; with HTTPS overhead (a TLS handshake for most messages) roughly 2 MB. Everything is gzip JSON over HTTPS, POST https://approvalens.com/api/agent/v1/report, with the server's token. When offline the agent keeps messages for up to an hour and sends them later.

What it never does

  • No remote commands

    It never executes commands or shells out, and never runs code it receives. The only thing it reads from an answer is the HTTP status code: 2xx means delivered, 205 means “send your full lists again”.

  • No inbound connections

    It opens no listening socket. The only traffic is an outgoing HTTPS POST to the configured URL (plain HTTP is refused except to 127.0.0.1).

  • No file contents

    It never reads the content of your files. Web root checks compare modification time and size and send file names only. SSH keys are sent as fingerprints.

  • No changes

    It never kills, quarantines or deletes anything. It writes only its own state folder, /var/lib/approvalens-agent (unsent messages for at most an hour, what it last sent, the web root file list).

Security signals are heuristics

A match is a reason to look, not proof, and no match proves nothing: this is not an antivirus. The signals are: a process name or argument from the crypto-miner list (signatures.txt below); an executable under /tmp, /var/tmp or /dev/shm; an executable deleted after it started (outside the system folders, which upgrades do legitimately) or running from memory (memfd); one process using a whole core for 15 minutes (databases, compressors, compilers and backup tools excepted); and anything new in the lists above after the first 24 hours (the baseline). You can acknowledge a signal and it joins the baseline.

Permissions

By default it runs as its own user, approvalens-agent, with no capabilities at all. Linux then hides a few things from it, and your account says which: the executable path of other users' processes (it falls back to the command line), the owner of other users' listening sockets, users' crontabs and other users' SSH keys. Install with --root to include them: the unit then runs as root but keeps only CAP_DAC_READ_SEARCH and CAP_SYS_PTRACE (to read those files and /proc/PID/exe and fd), on a read-only filesystem.

Resource limits and systemd hardening

Goal: under 1% of one core on average and under 30 MB of memory. The unit enforces hard caps and locks the process down:

[Unit]
Description=Approvalens server agent (reads system statistics, reports over HTTPS)
Documentation=https://approvalens.com/agent
After=network-online.target
Wants=network-online.target

[Service]
Type=simple
User=approvalens-agent
ExecStart=/usr/local/bin/approvalens-agent run --config /etc/approvalens-agent.yaml
Restart=always
RestartSec=15
StateDirectory=approvalens-agent
StateDirectoryMode=0700
# Hardening: everything read-only except its own state directory, no privilege gain.
NoNewPrivileges=yes
ProtectSystem=strict
ProtectHome=read-only
ReadOnlyPaths=/
CapabilityBoundingSet=
AmbientCapabilities=
ProtectKernelTunables=yes
ProtectKernelModules=yes
ProtectKernelLogs=yes
ProtectControlGroups=yes
ProtectClock=yes
ProtectHostname=yes
RestrictAddressFamilies=AF_INET AF_INET6 AF_UNIX
RestrictNamespaces=yes
RestrictRealtime=yes
RestrictSUIDSGID=yes
LockPersonality=yes
MemoryDenyWriteExecute=yes
SystemCallArchitectures=native
SystemCallFilter=@system-service
UMask=0077
# Footprint caps
CPUQuota=5%
MemoryMax=64M
Nice=10
IOSchedulingClass=idle

[Install]
WantedBy=multi-user.target

Measured on a test machine: 0.12% of one core on average, 12 MB of memory.

Install

Create a server in your account to get its token, then on the server: /account

curl -fsSL https://approvalens.com/agent/install.sh | sudo bash -s -- --token <TOKEN>

See every step first, without changing anything:

curl -fsSL https://approvalens.com/agent/install.sh | sudo bash -s -- --token <TOKEN> --dry-run

Install a binary you built yourself:

curl -fsSLO https://approvalens.com/agent/install.sh
sudo bash install.sh --token <TOKEN> --binary ./approvalens-agent-linux-amd64

Print exactly what it would send, without sending:

sudo -u approvalens-agent approvalens-agent check --config /etc/approvalens-agent.yaml

Uninstall

Stops the service and removes the binary, the config, the state folder, the unit and the user:

curl -fsSL https://approvalens.com/agent/install.sh | sudo bash -s -- --uninstall

Build it yourself and compare

The builds are reproducible: the same Go version (the toolchain line in go.mod; Go downloads it by itself) produces byte-identical binaries on any machine. Build from the source tarball and compare the SHA-256 with ours:

curl -fsSLO https://approvalens.com/agent/0.1.0/approvalens-agent-0.1.0-src.tar.gz
curl -fsSLO https://approvalens.com/agent/0.1.0/SHA256SUMS
sha256sum --ignore-missing -c SHA256SUMS
tar xzf approvalens-agent-0.1.0-src.tar.gz && cd approvalens-agent-0.1.0
go test ./...
for arch in amd64 arm64; do
  CGO_ENABLED=0 GOOS=linux GOARCH=$arch go build -trimpath -buildvcs=false \
    -ldflags "-s -w -buildid= -X main.version=0.1.0" -o approvalens-agent-linux-$arch .
done
sha256sum approvalens-agent-linux-amd64 approvalens-agent-linux-arm64
grep linux ../SHA256SUMS

The sums must equal the lines in SHA256SUMS. If they do not, do not install it, and tell us.

Downloads · version 0.1.0

8e6ab0397de8706f17304fe0392a284e5192a2b9886efc3315f48dff645b8539  approvalens-agent-linux-amd64
5d7ba23c0572174cc0135bd1fe9c7e0d9b04809ea7c462129bd88f089f3a0d49  approvalens-agent-linux-arm64
0b0a10618fa7916d47357cacee3091bb704f44a717f9b4139292b6ff176aa1db  approvalens-agent-0.1.0-src.tar.gz

Source code

The main files of version 0.1.0, exactly as built. The tarball above also has the tests and their fixtures.

main.go225 linesraw
// Command approvalens-agent: the optional Approvalens server agent.
//
// It reads /proc, /sys, statfs and a few files (crontabs, authorized_keys, the setuid bits
// of binaries, the size and time of files in configured web roots), aggregates a sample
// taken every 60 seconds and posts a compact report every 5 minutes to approvalens.com.
// It never executes commands, never runs code it receives, opens no listening socket and
// changes nothing on the machine apart from its own state directory.
package main

import (
	"encoding/json"
	"flag"
	"fmt"
	"log"
	"os"
	"os/signal"
	"runtime/debug"
	"strings"
	"syscall"
	"time"
)

var version = "dev" // set by build.sh (-ldflags "-X main.version=...")

func usage() {
	fmt.Fprintf(os.Stderr, `approvalens-agent %s

Usage:
  approvalens-agent run     [--config FILE]   sample every 60 s, report every 5 min (the systemd service)
  approvalens-agent check   [--config FILE]   take one sample and print the report it would send (sends nothing)
  approvalens-agent version

Config: %s (token, url, web_roots, intervals). See https://approvalens.com/monitoring
`, version, DefaultConfigPath)
}

func main() {
	log.SetFlags(0)
	if len(os.Args) < 2 {
		usage()
		os.Exit(2)
	}
	cmd := os.Args[1]
	fs := flag.NewFlagSet(cmd, flag.ExitOnError)
	cfgPath := fs.String("config", envOr("APPROVALENS_AGENT_CONFIG", DefaultConfigPath), "config file")
	switch cmd {
	case "version", "--version", "-v":
		fmt.Println("approvalens-agent", version)
		return
	case "run":
		_ = fs.Parse(os.Args[2:])
		os.Exit(run(*cfgPath))
	case "check":
		_ = fs.Parse(os.Args[2:])
		os.Exit(check(*cfgPath))
	case "help", "-h", "--help":
		usage()
	default:
		usage()
		os.Exit(2)
	}
}

func envOr(k, d string) string {
	if v := os.Getenv(k); v != "" {
		return v
	}
	return d
}

// Memory stays small: a soft limit well under the unit's MemoryMax and an eager GC.
func tuneRuntime() {
	debug.SetGCPercent(50)
	debug.SetMemoryLimit(24 << 20)
}

func run(path string) int {
	tuneRuntime()
	cfg, err := LoadConfig(path)
	if err != nil {
		log.Printf("config: %v", err)
		return 1
	}
	c := NewCollector(cfg, version)
	o := NewOutbox(c, true)
	s := NewSender(cfg, version)
	s.OnResync = o.Resync
	s.LoadSpool()
	c.LoadWebState()
	log.Printf("approvalens-agent %s: sampling every %s; heartbeat every %s, rollup every %s, events at once (min gap %s) to %s (uid %d)",
		version, cfg.SampleInterval, cfg.HeartbeatInterval, cfg.RollupInterval, cfg.EventGap, cfg.URL, os.Geteuid())
	if sk := c.Skipped(); len(sk) > 0 {
		log.Printf("checks skipped with these permissions: %s", strings.Join(sk, ", "))
	}

	slow := make(chan string, 4)
	go func() { // setuid and web root scans: slow file walks, off the sampling loop
		for job := range slow {
			switch job {
			case "setuid":
				c.ScanSetuidNow()
			case "web":
				c.ScanWebNow(true)
			}
		}
	}()
	queue := func(job string) {
		select {
		case slow <- job:
		default:
		}
	}

	stop := make(chan os.Signal, 1)
	signal.Notify(stop, syscall.SIGINT, syscall.SIGTERM)
	tick := time.NewTicker(cfg.SampleInterval)
	defer tick.Stop()

	start := time.Now()
	c.Sample(start)
	o.RefreshInventory()
	s.Enqueue(o.Heartbeat(start.UTC(), s.Pending())) // the account shows "connected" right away
	if o.NeedFull() {
		s.Enqueue(o.Full(start.UTC()))
	}
	nextHB := start.Add(cfg.HeartbeatInterval)
	nextRollup := start.Add(time.Minute) // the first chart points soon after install, then every 15 min
	nextInv := start.Add(cfg.HeartbeatInterval)
	nextSetuid := start.Add(2 * time.Minute)
	nextWeb := start
	lastStats := start
	for {
		select {
		case <-stop:
			msgs, bytes := s.Sent()
			log.Printf("stopping (%d messages, %d bytes sent since start)", msgs, bytes)
			return 0
		case now := <-tick.C:
			c.Sample(now)
			if !now.Before(nextSetuid) {
				nextSetuid = now.Add(cfg.SetuidInterval)
				queue("setuid")
			}
			if len(cfg.WebRoots) > 0 && !now.Before(nextWeb) {
				nextWeb = now.Add(cfg.WebRootInterval)
				queue("web")
			}
			if !now.Before(nextInv) { // crontabs and SSH keys: read every heartbeat interval
				nextInv = now.Add(cfg.HeartbeatInterval)
				o.RefreshInventory()
			}
			utc := now.UTC()
			if o.NeedFull() {
				o.RefreshInventory()
				s.Enqueue(o.Full(utc))
			}
			if m := o.Event(utc, cfg.EventGap); m != nil {
				s.Enqueue(*m)
			}
			if !now.Before(nextRollup) {
				nextRollup = now.Add(cfg.RollupInterval)
				s.Enqueue(o.Rollup(utc))
			}
			if !now.Before(nextHB) {
				nextHB = now.Add(cfg.HeartbeatInterval)
				s.Enqueue(o.Heartbeat(utc, s.Pending()))
			}
			if now.Sub(lastStats) >= 24*time.Hour {
				msgs, bytes := s.Sent()
				log.Printf("sent %d messages, %d bytes since start", msgs, bytes)
				lastStats = now
			}
		}
	}
}

// check prints the messages the agent would send now. Nothing is sent and no state is written.
func check(path string) int {
	cfg, err := LoadConfig(path)
	if err != nil {
		// Transparency first: show what would be collected even without a valid config.
		fmt.Fprintf(os.Stderr, "config: %v (continuing with defaults; nothing is sent by `check` anyway)\n", err)
		cfg = DefaultConfig()
		cfg.path = path
	}
	c := NewCollector(cfg, version)
	c.LoadWebState()
	o := NewOutbox(c, false)
	c.Sample(time.Now())
	time.Sleep(2 * time.Second) // a second sample turns counters into rates
	c.Sample(time.Now())
	c.ScanSetuidNow()
	c.ScanWebNow(false)
	o.RefreshInventory()
	now := time.Now().UTC()
	msgs := []any{o.Heartbeat(now, 0), o.Full(now), o.Rollup(now)}
	if m := o.Event(now, 0); m != nil {
		msgs = append(msgs, *m)
	}
	total := 0
	for _, m := range msgs {
		b, _ := json.MarshalIndent(m, "", "  ")
		gz, _ := Encode(m)
		total += len(gz)
		fmt.Println(string(b))
		fmt.Fprintf(os.Stderr, "-- %s: %d bytes gzip (%d bytes JSON)\n", m.(interface{ kind() string }).kind(), len(gz), len(b))
	}
	tok := cfg.Token
	if len(tok) > 9 {
		tok = tok[:9] + "…"
	}
	fmt.Fprintf(os.Stderr, "\nWould POST these to %s with token %s. Schedule: heartbeat every %s, rollup every %s, events at once (at most one per %s), the full inventory only after install or when the server asks for a resync.\n",
		cfg.Endpoint(), tok, cfg.HeartbeatInterval, cfg.RollupInterval, cfg.EventGap)
	fmt.Fprintf(os.Stderr, "Running as uid %d. Checks skipped with these permissions: %s\n", os.Geteuid(), joinOr(c.Skipped(), "none"))
	fmt.Fprintln(os.Stderr, "Nothing was sent. The agent never executes commands and has no listening socket.")
	return 0
}

func joinOr(xs []string, d string) string {
	if len(xs) == 0 {
		return d
	}
	return strings.Join(xs, ", ")
}
config.go237 linesraw
package main

import (
	"bufio"
	"bytes"
	"errors"
	"fmt"
	"net"
	"net/url"
	"os"
	"strconv"
	"strings"
	"time"
)

const DefaultConfigPath = "/etc/approvalens-agent.yaml"

// Config is /etc/approvalens-agent.yaml. Only a small YAML subset is read: `key: value`
// lines and lists written as `- item` lines under a key (or `[a, b]`).
type Config struct {
	Token             string
	URL               string
	WebRoots          []string
	WebRootExclude    []string
	WebRootMaxFiles   int
	SampleInterval    time.Duration
	HeartbeatInterval time.Duration
	RollupInterval    time.Duration
	EventGap          time.Duration
	SetuidInterval    time.Duration
	WebRootInterval   time.Duration
	StateDir          string
	path              string
}

func DefaultConfig() Config {
	return Config{
		URL:               "https://approvalens.com",
		WebRootMaxFiles:   20000,
		SampleInterval:    60 * time.Second,
		HeartbeatInterval: 5 * time.Minute,
		RollupInterval:    15 * time.Minute,
		EventGap:          60 * time.Second,
		SetuidInterval:    6 * time.Hour,
		WebRootInterval:   15 * time.Minute,
		StateDir:          "/var/lib/approvalens-agent",
	}
}

func LoadConfig(path string) (Config, error) {
	c := DefaultConfig()
	c.path = path
	b, err := os.ReadFile(path)
	if err != nil {
		return c, err
	}
	if err := c.parse(b); err != nil {
		return c, fmt.Errorf("%s: %w", path, err)
	}
	if env := os.Getenv("APPROVALENS_AGENT_TOKEN"); env != "" {
		c.Token = env
	}
	return c, c.Validate()
}

func (c *Config) parse(b []byte) error {
	sc := bufio.NewScanner(bytes.NewReader(b))
	listKey := ""
	n := 0
	for sc.Scan() {
		n++
		raw := sc.Text()
		line := strings.TrimSpace(stripComment(raw))
		if line == "" {
			continue
		}
		if strings.HasPrefix(line, "- ") || line == "-" {
			if listKey == "" {
				return fmt.Errorf("line %d: list item without a key", n)
			}
			if err := c.set(listKey, unquote(strings.TrimSpace(strings.TrimPrefix(line, "-"))), true); err != nil {
				return fmt.Errorf("line %d: %w", n, err)
			}
			continue
		}
		k, v, ok := strings.Cut(line, ":")
		if !ok {
			return fmt.Errorf("line %d: expected `key: value`", n)
		}
		k, v = strings.TrimSpace(k), strings.TrimSpace(v)
		listKey = ""
		if v == "" {
			listKey = k
			continue
		}
		if strings.HasPrefix(v, "[") && strings.HasSuffix(v, "]") {
			for _, item := range strings.Split(strings.Trim(v, "[]"), ",") {
				if item = unquote(strings.TrimSpace(item)); item != "" {
					if err := c.set(k, item, true); err != nil {
						return fmt.Errorf("line %d: %w", n, err)
					}
				}
			}
			continue
		}
		if err := c.set(k, unquote(v), false); err != nil {
			return fmt.Errorf("line %d: %w", n, err)
		}
	}
	return sc.Err()
}

func stripComment(s string) string {
	inQ := byte(0)
	for i := 0; i < len(s); i++ {
		ch := s[i]
		switch {
		case inQ != 0 && ch == inQ:
			inQ = 0
		case inQ == 0 && (ch == '"' || ch == '\''):
			inQ = ch
		case inQ == 0 && ch == '#' && (i == 0 || s[i-1] == ' ' || s[i-1] == '\t'):
			return s[:i]
		}
	}
	return s
}

func unquote(s string) string {
	if len(s) >= 2 && (s[0] == '"' && s[len(s)-1] == '"' || s[0] == '\'' && s[len(s)-1] == '\'') {
		return s[1 : len(s)-1]
	}
	return s
}

func (c *Config) set(k, v string, list bool) error {
	switch k {
	case "token":
		c.Token = v
	case "url":
		c.URL = strings.TrimRight(v, "/")
	case "web_roots":
		c.WebRoots = append(c.WebRoots, v)
	case "web_root_exclude":
		c.WebRootExclude = append(c.WebRootExclude, v)
	case "web_root_max_files":
		n, err := strconv.Atoi(v)
		if err != nil || n < 1 {
			return fmt.Errorf("web_root_max_files: %q", v)
		}
		c.WebRootMaxFiles = n
	case "sample_interval", "heartbeat_interval", "rollup_interval", "event_min_gap", "setuid_interval", "web_root_interval":
		d, err := parseDuration(v)
		if err != nil {
			return fmt.Errorf("%s: %w", k, err)
		}
		switch k {
		case "sample_interval":
			c.SampleInterval = d
		case "heartbeat_interval":
			c.HeartbeatInterval = d
		case "rollup_interval":
			c.RollupInterval = d
		case "event_min_gap":
			c.EventGap = d
		case "setuid_interval":
			c.SetuidInterval = d
		case "web_root_interval":
			c.WebRootInterval = d
		}
	case "state_dir":
		c.StateDir = v
	default:
		return fmt.Errorf("unknown key %q", k)
	}
	_ = list
	return nil
}

// parseDuration accepts Go durations ("90s", "5m", "6h") or plain seconds.
func parseDuration(v string) (time.Duration, error) {
	if n, err := strconv.Atoi(v); err == nil {
		return time.Duration(n) * time.Second, nil
	}
	return time.ParseDuration(v)
}

func (c Config) Validate() error {
	if c.Token == "" {
		return errors.New("token is missing")
	}
	for _, r := range c.Token {
		if !(r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z' || r >= '0' && r <= '9' || r == '_' || r == '-') {
			return errors.New("token has unexpected characters")
		}
	}
	u, err := url.Parse(c.URL)
	if err != nil || u.Host == "" {
		return fmt.Errorf("url %q is not valid", c.URL)
	}
	// Plain HTTP only to this machine (local testing); anything else must be HTTPS.
	if u.Scheme != "https" {
		host := u.Hostname()
		ip := net.ParseIP(host)
		if u.Scheme != "http" || !(host == "localhost" || (ip != nil && ip.IsLoopback())) {
			return fmt.Errorf("url must be https:// (got %q)", c.URL)
		}
	}
	if c.SampleInterval < 5*time.Second || c.SampleInterval > 10*time.Minute {
		return errors.New("sample_interval must be between 5s and 10m")
	}
	if c.HeartbeatInterval < 30*time.Second || c.HeartbeatInterval > 10*time.Minute || c.HeartbeatInterval < c.SampleInterval {
		return errors.New("heartbeat_interval must be between 30s and 10m, and not shorter than sample_interval")
	}
	if c.RollupInterval < c.HeartbeatInterval || c.RollupInterval > time.Hour {
		return errors.New("rollup_interval must be between heartbeat_interval and 1h")
	}
	if c.EventGap < 10*time.Second || c.EventGap > 10*time.Minute {
		return errors.New("event_min_gap must be between 10s and 10m")
	}
	if c.SetuidInterval < time.Minute {
		return errors.New("setuid_interval must be at least 1m")
	}
	if c.WebRootInterval < time.Minute {
		return errors.New("web_root_interval must be at least 1m")
	}
	for _, r := range c.WebRoots {
		if !strings.HasPrefix(r, "/") {
			return fmt.Errorf("web_roots: %q is not an absolute path", r)
		}
	}
	return nil
}

// Endpoint the reports go to.
func (c Config) Endpoint() string { return c.URL + "/api/agent/v1/report" }
collector.go609 linesraw
package main

import (
	"encoding/json"
	"os"
	"path/filepath"
	"runtime"
	"sort"
	"strconv"
	"strings"
	"sync"
	"time"
)

const (
	hogPercent  = 90.0             // one process at this share of a core ...
	hogDuration = 15 * time.Minute // ... for this long, without a break
	fillWindow  = 6 * 3600         // disk fill projection looks at the last 6 hours
)

// ---------------------------------------------------------------------------
// Data shared by the messages (messages.go). Field names are short: they go over the network.
// ---------------------------------------------------------------------------

type Minute struct {
	T      int64   `json:"t"`       // start of the minute, unix seconds (UTC)
	CPU    float64 `json:"cpu"`     // % of all cores, average
	CPUMax float64 `json:"cpu_max"` // highest sample in the minute
	Load1  float64 `json:"load1"`
	Mem    float64 `json:"mem"` // % used (MemAvailable based)
	MemMax float64 `json:"mem_max"`
	Swap   float64 `json:"swap"`    // % of swap used
	SwapIO float64 `json:"swap_io"` // pages swapped in + out per second
	RX     float64 `json:"rx"`      // bytes per second
	TX     float64 `json:"tx"`
	Disk   float64 `json:"disk"` // fullest mount, % used
	N      int     `json:"n"`    // samples in the minute
}

type DiskOut struct {
	Mount     string   `json:"mount"`
	FS        string   `json:"fs"`
	Total     uint64   `json:"total"`
	Used      uint64   `json:"used"`
	Pct       float64  `json:"pct"`
	InodesPct float64  `json:"inodes_pct"`
	RateBPH   *float64 `json:"rate_bph,omitempty"`  // bytes per hour over the last 6 h (least squares)
	FullInH   *float64 `json:"full_in_h,omitempty"` // hours until full at that rate
}

type PortOut struct {
	Proto string `json:"proto"`
	Addr  string `json:"addr"`
	Port  int    `json:"port"`
	Local bool   `json:"local"` // bound to loopback only
	PID   int    `json:"pid,omitempty"`
	Name  string `json:"name,omitempty"`
	User  string `json:"user,omitempty"`
}

func (p PortOut) Key() string {
	k := p.Proto + "/" + strconv.Itoa(p.Port)
	if p.Local {
		k += "/local"
	}
	return k
}

type MemOut struct {
	Total     uint64  `json:"total"`
	Available uint64  `json:"available"`
	UsedPct   float64 `json:"used_pct"`
	SwapTotal uint64  `json:"swap_total"`
	SwapUsed  uint64  `json:"swap_used"`
	SwapPct   float64 `json:"swap_pct"`
}

type Current struct {
	At     int64      `json:"at"`
	CPU    float64    `json:"cpu"`
	Load   [3]float64 `json:"load"`
	Mem    MemOut     `json:"mem"`
	Disks  []DiskOut  `json:"disks"`
	RX     float64    `json:"rx"`
	TX     float64    `json:"tx"`
	Procs  int        `json:"procs"`
	TopCPU []Proc     `json:"top_cpu"`
	TopMem []Proc     `json:"top_mem"`
	Ports  []PortOut  `json:"ports"`
}

type Signal struct {
	Kind   string         `json:"kind"` // miner | tmp_exe | deleted_exe | cpu_hog
	Key    string         `json:"key"`  // stable across samples and restarts of the process
	At     int64          `json:"at"`   // first seen in this report window
	Detail map[string]any `json:"detail"`
}

type WebRootOut struct {
	Roots    []string    `json:"roots"`
	Files    int         `json:"files"`
	LastScan int64       `json:"last_scan,omitempty"`
	Baseline bool        `json:"baseline,omitempty"` // first scan: nothing to compare with yet
	Changes  *WebChanges `json:"changes,omitempty"`
}

type AgentInfo struct {
	Version   string   `json:"version"`
	Root      bool     `json:"root"`
	UID       int      `json:"uid"`
	Skipped   []string `json:"skipped"`
	SampleS   int      `json:"sample_s"`
	ReportS   int      `json:"report_s"`
	StartedAt int64    `json:"started_at"`
}

type HostInfo struct {
	Hostname string `json:"hostname"`
	OS       string `json:"os"`
	Kernel   string `json:"kernel"`
	Arch     string `json:"arch"`
	Cores    int    `json:"cores"`
	BootTime int64  `json:"boot_time"`
	UptimeS  int64  `json:"uptime_s"`
}

// ---------------------------------------------------------------------------
// Collector
// ---------------------------------------------------------------------------

type minuteAgg struct {
	n                                                          int
	cpu, cpuMax, load, mem, memMax, swap, swapIO, rx, tx, disk float64
}

type Collector struct {
	cfg     Config
	version string
	proc    string // /proc
	host    string // / (root of the files read)
	users   *UserCache
	procs   *ProcSampler
	started time.Time

	prevT       time.Time
	prevCPU     CPUTimes
	prevNet     NetCounters
	prevSwapIn  uint64
	prevSwapOut uint64
	havePrev    bool

	fill      map[string]*FillTracker
	minutes   map[int64]*minuteAgg
	sampleSig map[string]Signal    // process signals seen in the latest sample
	hog       map[string]time.Time // pid:start -> first sample over hogPercent
	rules     *Rules
	pending   []Event // tripped since the last event message
	lastSwap  float64

	portInodes map[uint64]bool
	portOwners map[uint64]Proc

	current  Current
	hostInfo HostInfo
	samples  int
	errors   int
	seq      int64

	mu       sync.Mutex // guards the scan results below (filled by the slow scans goroutine)
	setuid   []string   // nil until the first scan
	setuidAt int64
	web      *WebRootOut
	webPrev  map[string]string
	webPend  *WebChanges

	cronUnreadable []string
	sshUnreadable  []string
}

func NewCollector(cfg Config, version string) *Collector {
	users := NewUserCache("/etc/passwd")
	c := &Collector{
		cfg: cfg, version: version, proc: "/proc", host: "/", users: users,
		procs: NewProcSampler("/proc", users), started: time.Now(),
		fill: map[string]*FillTracker{}, minutes: map[int64]*minuteAgg{}, sampleSig: map[string]Signal{},
		hog: map[string]time.Time{}, portInodes: map[uint64]bool{}, portOwners: map[uint64]Proc{}, rules: NewRules(),
	}
	return c
}

func (c *Collector) read(name string) ([]byte, error) {
	return os.ReadFile(filepath.Join(c.proc, name))
}

// Sample takes one reading of everything that is sampled every interval.
func (c *Collector) Sample(now time.Time) {
	c.samples++
	ok := true
	var cur Current
	cur.At = now.Unix()

	b, err := c.read("stat")
	ps, perr := ParseProcStat(b)
	if err != nil || perr != nil {
		ok = false
	}
	b, _ = c.read("meminfo")
	mem, merr := ParseMemInfo(b)
	if merr != nil {
		ok = false
	}
	b, _ = c.read("vmstat")
	swIn, swOut := ParseVMStat(b)
	b, _ = c.read("net/dev")
	netc := ParseNetDev(b)
	b, _ = c.read("loadavg")
	cur.Load, _ = ParseLoadAvg(b)
	b, _ = c.read("uptime")
	up, _ := ParseUptime(b)

	elapsed := now.Sub(c.prevT).Seconds()
	var swapIO float64
	if c.havePrev && elapsed > 0 {
		cur.CPU = round1(CPUPercent(c.prevCPU, ps.CPU))
		if netc.RX >= c.prevNet.RX {
			cur.RX = float64(netc.RX-c.prevNet.RX) / elapsed
		}
		if netc.TX >= c.prevNet.TX {
			cur.TX = float64(netc.TX-c.prevNet.TX) / elapsed
		}
		if swIn >= c.prevSwapIn && swOut >= c.prevSwapOut {
			swapIO = float64(swIn-c.prevSwapIn+swOut-c.prevSwapOut) / elapsed
		}
	}
	cur.RX, cur.TX = float64(int64(cur.RX)), float64(int64(cur.TX))
	c.prevCPU, c.prevNet, c.prevSwapIn, c.prevSwapOut, c.prevT, c.havePrev = ps.CPU, netc, swIn, swOut, now, true

	cur.Mem = MemOut{Total: mem.Total, Available: mem.Available, UsedPct: round1(mem.UsedPct()), SwapTotal: mem.SwapTotal,
		SwapUsed: mem.SwapUsed(), SwapPct: round1(mem.SwapPct())}

	// Disks
	b, _ = c.read("self/mountinfo")
	maxDisk := 0.0
	seenMounts := map[string]bool{}
	for _, m := range ParseMountInfo(b) {
		d, err := statDisk(m.Point)
		if err != nil || d.Total == 0 {
			continue
		}
		seenMounts[m.Point] = true
		ft := c.fill[m.Point]
		if ft == nil {
			ft = NewFillTracker(fillWindow)
			c.fill[m.Point] = ft
		}
		ft.Add(now.Unix(), d.Used)
		o := DiskOut{Mount: m.Point, FS: m.FSType, Total: d.Total, Used: d.Used, Pct: round1(d.UsedPct), InodesPct: round1(d.InodesPct)}
		if r, ok := ft.Rate(); ok {
			r = float64(int64(r))
			o.RateBPH = &r
			if h, ok := ft.FullInHours(d.Avail); ok {
				h = round1(h)
				o.FullInH = &h
			}
		}
		if o.Pct > maxDisk {
			maxDisk = o.Pct
		}
		cur.Disks = append(cur.Disks, o)
	}
	for k := range c.fill {
		if !seenMounts[k] {
			delete(c.fill, k)
		}
	}

	// Processes
	procs := c.procs.Sample(now)
	cur.Procs = len(procs)
	cur.TopCPU = TopBy(procs, 10, func(p Proc) float64 { return p.CPU })
	cur.TopMem = TopBy(procs, 10, func(p Proc) float64 { return float64(p.RSS) })
	c.processSignals(now, procs)

	// Listening sockets
	cur.Ports = c.ports(procs)

	c.current = cur
	c.hostInfo = c.readHost(ps, up)
	c.lastSwap = swapIO

	// Local rules: what trips now goes out in the next event message (and only then).
	evs, _ := c.rules.Update(now, Reading{Mem: cur.Mem.UsedPct, SwapIO: swapIO, SwapPct: cur.Mem.SwapPct, Load1: cur.Load[0],
		Cores: ps.Cores, Disks: cur.Disks, TopCPU: cur.TopCPU, TopMem: cur.TopMem, Signals: c.sampleSig})
	c.pending = append(c.pending, evs...)
	if len(c.pending) > 100 {
		c.pending = c.pending[len(c.pending)-100:]
	}

	if !ok {
		c.errors++
		return
	}
	t := now.Unix() - now.Unix()%60
	a := c.minutes[t]
	if a == nil {
		a = &minuteAgg{}
		c.minutes[t] = a
	}
	a.n++
	a.cpu += cur.CPU
	a.cpuMax = maxF(a.cpuMax, cur.CPU)
	a.load += cur.Load[0]
	a.mem += cur.Mem.UsedPct
	a.memMax = maxF(a.memMax, cur.Mem.UsedPct)
	a.swap += cur.Mem.SwapPct
	a.swapIO += swapIO
	a.rx += cur.RX
	a.tx += cur.TX
	a.disk = maxF(a.disk, maxDisk)
}

func (c *Collector) readHost(ps ProcStat, up float64) HostInfo {
	h := HostInfo{Arch: runtime.GOARCH, Cores: ps.Cores, BootTime: ps.BootTime, UptimeS: int64(up)}
	if b, err := c.read("sys/kernel/hostname"); err == nil {
		h.Hostname = strings.TrimSpace(string(b))
	}
	if b, err := c.read("sys/kernel/osrelease"); err == nil {
		h.Kernel = strings.TrimSpace(string(b))
	}
	if b, err := os.ReadFile(filepath.Join(c.host, "etc/os-release")); err == nil {
		h.OS = ParseOSRelease(b)
	} else if b, err := os.ReadFile(filepath.Join(c.host, "usr/lib/os-release")); err == nil {
		h.OS = ParseOSRelease(b)
	}
	return h
}

func (c *Collector) processSignals(now time.Time, procs []Proc) {
	c.sampleSig = map[string]Signal{}
	live := map[string]bool{}
	for _, p := range procs {
		if p.kernel || p.PID == os.Getpid() {
			continue
		}
		// Crypto miner by name or arguments
		if why := sigs.MinerMatch(p.Name, p.argv0, p.Cmd); why != "" {
			c.addSignal(now, "miner", "miner:"+strings.ToLower(p.Name)+":"+exeOrArgv0(p), p, map[string]any{"match": why})
		}
		// Executable in a temporary directory, or deleted after start
		exe := strings.TrimSuffix(p.Exe, " (deleted)")
		switch {
		case p.exeKnown && InTmp(exe):
			c.addSignal(now, "tmp_exe", "tmp_exe:"+exe, p, map[string]any{"source": "exe"})
		case !p.exeKnown && strings.HasPrefix(p.argv0, "/") && InTmp(p.argv0):
			c.addSignal(now, "tmp_exe", "tmp_exe:"+p.argv0, p, map[string]any{"source": "argv0"})
		}
		if p.exeKnown && SuspiciousDeleted(p.Exe) {
			c.addSignal(now, "deleted_exe", "deleted_exe:"+exe, p, map[string]any{"source": "exe"})
		}
		// Sustained CPU by one process
		id := strconv.Itoa(p.PID) + ":" + strconv.FormatUint(p.startTime, 10)
		live[id] = true
		if p.CPU >= hogPercent {
			first, seen := c.hog[id]
			if !seen {
				c.hog[id] = now
				first = now
			}
			if now.Sub(first) >= hogDuration && !sigs.Busy[strings.ToLower(p.Name)] {
				c.addSignal(now, "cpu_hog", "cpu_hog:"+strings.ToLower(p.Name), p,
					map[string]any{"minutes": int(now.Sub(first).Minutes()), "since": first.Unix()})
			}
		} else {
			delete(c.hog, id)
		}
	}
	for id := range c.hog {
		if !live[id] {
			delete(c.hog, id)
		}
	}
}

func exeOrArgv0(p Proc) string {
	if p.exeKnown {
		return strings.TrimSuffix(p.Exe, " (deleted)")
	}
	return p.argv0
}

func (c *Collector) addSignal(now time.Time, kind, key string, p Proc, detail map[string]any) {
	key = truncate(key, 300)
	detail["process"] = p.Reported()
	c.sampleSig[key] = Signal{Kind: kind, Key: key, At: now.Unix(), Detail: detail}
}

// localPortRange: UDP sockets bound to a port in this range are clients (resolvers,
// time sync), not services.
func (c *Collector) localPortRange() (int, int) {
	b, err := c.read("sys/net/ipv4/ip_local_port_range")
	if err == nil {
		f := strings.Fields(string(b))
		if len(f) == 2 {
			lo, e1 := strconv.Atoi(f[0])
			hi, e2 := strconv.Atoi(f[1])
			if e1 == nil && e2 == nil {
				return lo, hi
			}
		}
	}
	return 32768, 60999
}

func (c *Collector) ports(procs []Proc) []PortOut {
	var socks []Socket
	for _, f := range [][2]string{{"net/tcp", "tcp"}, {"net/tcp6", "tcp"}, {"net/udp", "udp"}, {"net/udp6", "udp"}} {
		b, err := c.read(f[0])
		if err != nil {
			continue
		}
		s, _ := ParseNetSockets(b, f[1])
		socks = append(socks, s...)
	}
	lo, hi := c.localPortRange()
	inodes := map[uint64]bool{}
	var kept []Socket
	for _, s := range socks {
		if s.Proto == "udp" && s.Port >= lo && s.Port <= hi {
			continue
		}
		kept = append(kept, s)
		if s.Inode != 0 {
			inodes[s.Inode] = true
		}
	}
	if !sameSet(inodes, c.portInodes) {
		c.portOwners = SocketOwners(c.proc, inodes, procs)
		c.portInodes = inodes
	}
	seen := map[string]bool{}
	var out []PortOut
	for _, s := range kept {
		po := PortOut{Proto: s.Proto, Addr: s.Addr, Port: s.Port, Local: s.Local()}
		if p, ok := c.portOwners[s.Inode]; ok {
			po.PID, po.Name, po.User = p.PID, p.Name, p.User
		} else {
			po.User = c.users.Name(s.UID)
		}
		k := po.Key() + "|" + po.Addr
		if seen[k] {
			continue
		}
		seen[k] = true
		out = append(out, po)
	}
	sort.Slice(out, func(i, j int) bool {
		if out[i].Proto != out[j].Proto {
			return out[i].Proto < out[j].Proto
		}
		if out[i].Port != out[j].Port {
			return out[i].Port < out[j].Port
		}
		return out[i].Addr < out[j].Addr
	})
	if len(out) > 200 {
		out = out[:200]
	}
	return out
}

func sameSet(a, b map[uint64]bool) bool {
	if len(a) != len(b) {
		return false
	}
	for k := range a {
		if !b[k] {
			return false
		}
	}
	return true
}

// ---------------------------------------------------------------------------
// Slow scans (setuid every 6 h, web roots every 15 min), on their own goroutine
// ---------------------------------------------------------------------------

func (c *Collector) ScanSetuidNow() {
	items, truncated := ScanSetuid(setuidRoots, 200000)
	if len(items) > 500 {
		items = items[:500]
	}
	_ = truncated
	if items == nil {
		items = []string{}
	}
	c.mu.Lock()
	c.setuid, c.setuidAt = items, time.Now().Unix()
	c.mu.Unlock()
}

func (c *Collector) webStatePath() string { return filepath.Join(c.cfg.StateDir, "webroot.json") }

// LoadWebState reads the previous web root scan (kept across restarts).
func (c *Collector) LoadWebState() {
	b, err := os.ReadFile(c.webStatePath())
	if err != nil {
		return
	}
	var m map[string]string
	if json.Unmarshal(b, &m) == nil {
		c.mu.Lock()
		c.webPrev = m
		c.mu.Unlock()
	}
}

// ScanWebNow compares the web roots with the previous scan. persist=false (the `check`
// command) leaves the saved state alone.
func (c *Collector) ScanWebNow(persist bool) {
	if len(c.cfg.WebRoots) == 0 {
		return
	}
	files, truncated := ScanWebRoots(c.cfg.WebRoots, c.cfg.WebRootExclude, c.cfg.WebRootMaxFiles)
	c.mu.Lock()
	defer c.mu.Unlock()
	out := &WebRootOut{Roots: c.cfg.WebRoots, Files: len(files), LastScan: time.Now().Unix()}
	if c.webPrev == nil {
		out.Baseline = true
	} else {
		ch := DiffWebRoots(c.webPrev, files, 50)
		ch.Truncated = truncated
		if !ch.Empty() {
			c.webPend = mergeWeb(c.webPend, &ch)
		}
	}
	c.web = out
	c.webPrev = files
	if persist {
		if b, err := json.Marshal(files); err == nil {
			_ = writeFileAtomic(c.webStatePath(), b, 0o600)
		}
	}
}

func mergeWeb(a, b *WebChanges) *WebChanges {
	if a == nil {
		return b
	}
	m := *a
	m.Added = capList(dedupSorted(sortedCopy(append(m.Added, b.Added...))), 50)
	m.Changed = capList(dedupSorted(sortedCopy(append(m.Changed, b.Changed...))), 50)
	m.Removed = capList(dedupSorted(sortedCopy(append(m.Removed, b.Removed...))), 50)
	m.NAdded += b.NAdded
	m.NChanged += b.NChanged
	m.NRemoved += b.NRemoved
	m.Files = b.Files
	m.Truncated = m.Truncated || b.Truncated
	return &m
}

func capList(xs []string, n int) []string {
	if len(xs) > n {
		return xs[:n]
	}
	return xs
}

// ---------------------------------------------------------------------------
// Building the report
// ---------------------------------------------------------------------------

// Skipped names the checks this agent cannot do with its permissions.
func (c *Collector) Skipped() []string {
	var s []string
	if os.Geteuid() != 0 {
		s = append(s, "exe_paths_other_users", "port_owners_other_users")
	}
	if len(c.cronUnreadable) > 0 {
		s = append(s, "crontabs_unreadable")
	}
	if len(c.sshUnreadable) > 0 {
		s = append(s, "ssh_keys_unreadable")
	}
	if len(c.cfg.WebRoots) == 0 {
		s = append(s, "web_roots_not_configured")
	}
	return s
}

func maxF(a, b float64) float64 {
	if a > b {
		return a
	}
	return b
}

func round2(v float64) float64 { return float64(int64(v*100+0.5)) / 100 }

func writeFileAtomic(path string, b []byte, mode os.FileMode) error {
	if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil {
		return err
	}
	tmp := path + ".tmp"
	if err := os.WriteFile(tmp, b, mode); err != nil {
		return err
	}
	return os.Rename(tmp, path)
}
messages.go324 linesraw
package main

// What goes over the network. Sampling is local (every 60 s); the network sees four kinds
// of message, all gzip JSON posted to <url>/api/agent/v1/report:
//
//   heartbeat  every 5 min   a few numbers (CPU, memory, worst disk and its fill ETA, active rules)
//   rollup     every 15 min  per-minute averages / maxima since the last rollup, top processes,
//                            disks, network; once a day also a hash of each inventory list
//   event      at once       a local rule tripped (disk, inodes, memory, swap, load, a security
//                            signal) or an inventory list changed; at most one per minute
//   inventory  rarely        the full lists (listening ports, crontab lines, SSH key
//                            fingerprints, setuid files): after install, and when the server
//                            answers a daily hash with 205 (its copy drifted)

import (
	"crypto/sha256"
	"encoding/hex"
	"encoding/json"
	"os"
	"path/filepath"
	"sort"
	"strings"
	"sync/atomic"
	"time"
)

type Header struct {
	V       int    `json:"v"`
	Type    string `json:"type"`
	T       int64  `json:"t"` // unix seconds, UTC
	Seq     int64  `json:"seq"`
	Version string `json:"version"`
}

type HeartbeatMsg struct {
	Header
	UptimeS   int64    `json:"uptime_s"`
	CPU       float64  `json:"cpu"`
	Mem       float64  `json:"mem"`
	Disk      float64  `json:"disk"` // fullest filesystem, % used
	DiskMount string   `json:"disk_mount"`
	FullInH   *float64 `json:"full_in_h,omitempty"` // soonest projected full, hours
	Active    []string `json:"active"`              // rule keys tripped now
	Spooled   int      `json:"spooled"`
}

type RollupMsg struct {
	Header
	Agent   AgentInfo         `json:"agent"`
	Host    HostInfo          `json:"host"`
	Minutes []Minute          `json:"minutes"`
	Current Current           `json:"current"`
	Active  []string          `json:"active"`
	InvHash map[string]string `json:"inv_hash,omitempty"` // once a day
	WebRoot *WebRootOut       `json:"web_root,omitempty"`
}

type SetDiff struct {
	Added   []string `json:"added,omitempty"`
	Removed []string `json:"removed,omitempty"`
}

type EventMsg struct {
	Header
	Events  []Event            `json:"events,omitempty"`
	InvDiff map[string]SetDiff `json:"inv_diff,omitempty"`
	Ports   []PortOut          `json:"ports,omitempty"` // details of newly listening ports
	Web     *WebChanges        `json:"web,omitempty"`
}

type InventoryMsg struct {
	Header
	Reason    string    `json:"reason"` // install | resync
	Inventory Inv       `json:"inventory"`
	Ports     []PortOut `json:"ports"`
}

// Inv is the inventory as last sent. Setuid is nil until the first scan.
type Inv struct {
	Ports   []string `json:"ports"`
	Cron    []string `json:"cron"`
	SSHKeys []string `json:"ssh_keys"`
	Setuid  []string `json:"setuid"`
}

func (i *Inv) fields() map[string]*[]string {
	return map[string]*[]string{"ports": &i.Ports, "cron": &i.Cron, "ssh_keys": &i.SSHKeys, "setuid": &i.Setuid}
}

// HashItems is what the server recomputes over its copy of a list: sha256 of the sorted
// items joined by newlines, first 16 bytes in hex.
func HashItems(items []string) string {
	s := sortedCopy(items)
	sum := sha256.Sum256([]byte(strings.Join(s, "\n")))
	return hex.EncodeToString(sum[:16])
}

// DiffInv compares the inventory as sent with the current one. A list that is nil in
// cur (setuid before its first scan) is left out.
func DiffInv(sent, cur Inv) map[string]SetDiff {
	out := map[string]SetDiff{}
	sf, cf := sent.fields(), cur.fields()
	for name, cp := range cf {
		if *cp == nil {
			continue
		}
		a, r := DiffSets(*sf[name], *cp)
		if len(a)+len(r) > 0 {
			out[name] = SetDiff{Added: a, Removed: r}
		}
	}
	return out
}

// invState is kept in the state directory: what the server has, as far as the agent knows.
type invState struct {
	Sent        Inv   `json:"sent"`
	InstalledAt int64 `json:"installed_at"`
	HashDay     int64 `json:"hash_day"`
}

type Outbox struct {
	c        *Collector
	persist  bool
	state    invState
	haveFull bool // the server has a full inventory from us
	resync   atomic.Bool
	seq      int64
	lastEvt  time.Time
	cur      Inv
}

func (o *Outbox) path() string { return filepath.Join(o.c.cfg.StateDir, "inventory.json") }

func NewOutbox(c *Collector, persist bool) *Outbox {
	o := &Outbox{c: c, persist: persist}
	if persist {
		if b, err := os.ReadFile(o.path()); err == nil && json.Unmarshal(b, &o.state) == nil && o.state.InstalledAt > 0 {
			o.haveFull = true
		}
	}
	return o
}

func (o *Outbox) save() {
	if !o.persist {
		return
	}
	if b, err := json.Marshal(o.state); err == nil {
		_ = writeFileAtomic(o.path(), b, 0o600)
	}
}

func (o *Outbox) header(typ string, now time.Time) Header {
	o.seq++
	return Header{V: 1, Type: typ, T: now.Unix(), Seq: o.seq, Version: o.c.version}
}

// Resync is called when the server answered 205: the next message is the full inventory.
func (o *Outbox) Resync() { o.resync.Store(true) }

// RefreshInventory reads the lists that are cheap to read (ports from the last sample,
// crontabs and authorized_keys now) and takes the latest setuid scan.
func (o *Outbox) RefreshInventory() {
	c := o.c
	pk := map[string]bool{}
	for _, p := range c.current.Ports {
		pk[p.Key()] = true
	}
	ports := make([]string, 0, len(pk))
	for k := range pk {
		ports = append(ports, k)
	}
	sort.Strings(ports)
	cron, cu := ReadCrontabs(c.host)
	keys, ku := ReadSSHKeys(c.host, c.users.All())
	c.cronUnreadable, c.sshUnreadable = cu, ku
	c.mu.Lock()
	setuid := c.setuid
	c.mu.Unlock()
	o.cur = Inv{Ports: ports, Cron: nonNil(cron), SSHKeys: nonNil(keys), Setuid: setuid}
}

func nonNil(xs []string) []string {
	if xs == nil {
		return []string{}
	}
	return xs
}

func (o *Outbox) portDetails(keys []string) []PortOut {
	want := map[string]bool{}
	for _, k := range keys {
		want[k] = true
	}
	var out []PortOut
	for _, p := range o.c.current.Ports {
		if want[p.Key()] {
			out = append(out, p)
		}
	}
	return out
}

// NeedFull: after install (no saved state) or a 205 from the server.
func (o *Outbox) NeedFull() bool { return !o.haveFull || o.resync.Load() }

func (o *Outbox) Full(now time.Time) InventoryMsg {
	reason := "install"
	if o.haveFull {
		reason = "resync"
	}
	m := InventoryMsg{Header: o.header("inventory", now), Reason: reason, Inventory: o.cur, Ports: o.portDetails(o.cur.Ports)}
	if m.Inventory.Setuid == nil && o.state.Sent.Setuid != nil {
		m.Inventory.Setuid = o.state.Sent.Setuid
	}
	o.state.Sent = m.Inventory
	if o.state.InstalledAt == 0 {
		o.state.InstalledAt = now.Unix()
	}
	o.haveFull = true
	o.resync.Store(false)
	o.save()
	return m
}

func (o *Outbox) Heartbeat(now time.Time, spooled int) HeartbeatMsg {
	cur := o.c.current
	m := HeartbeatMsg{Header: o.header("heartbeat", now), UptimeS: o.c.hostInfo.UptimeS, CPU: cur.CPU, Mem: cur.Mem.UsedPct,
		Active: o.c.rules.Active(), Spooled: spooled}
	for _, d := range cur.Disks {
		if d.Pct >= m.Disk {
			m.Disk, m.DiskMount = d.Pct, d.Mount
		}
		if d.FullInH != nil && (m.FullInH == nil || *d.FullInH < *m.FullInH) {
			v := *d.FullInH
			m.FullInH = &v
		}
	}
	return m
}

// Rollup closes the chart window: per-minute aggregates since the last rollup.
func (o *Outbox) Rollup(now time.Time) RollupMsg {
	c := o.c
	m := RollupMsg{Header: o.header("rollup", now), Host: c.hostInfo, Active: c.rules.Active()}
	m.Agent = AgentInfo{Version: c.version, Root: os.Geteuid() == 0, UID: os.Geteuid(), SampleS: int(c.cfg.SampleInterval / time.Second),
		ReportS: int(c.cfg.RollupInterval / time.Second), StartedAt: c.started.Unix(), Skipped: nonNil(c.Skipped())}
	keys := make([]int64, 0, len(c.minutes))
	for t := range c.minutes {
		keys = append(keys, t)
	}
	sort.Slice(keys, func(i, j int) bool { return keys[i] < keys[j] })
	for _, t := range keys {
		a := c.minutes[t]
		n := float64(a.n)
		m.Minutes = append(m.Minutes, Minute{T: t, CPU: round1(a.cpu / n), CPUMax: round1(a.cpuMax), Load1: round2(a.load / n),
			Mem: round1(a.mem / n), MemMax: round1(a.memMax), Swap: round1(a.swap / n), SwapIO: round1(a.swapIO / n),
			RX: float64(int64(a.rx / n)), TX: float64(int64(a.tx / n)), Disk: round1(a.disk), N: a.n})
	}
	c.minutes = map[int64]*minuteAgg{}
	m.Current = c.current
	m.Current.Ports = nil // ports travel as inventory diffs
	c.mu.Lock()
	if c.web != nil {
		w := *c.web
		w.Changes = nil // changes travel in event messages
		m.WebRoot = &w
	}
	c.mu.Unlock()
	day := now.Unix() / 86400
	if o.haveFull && day != o.state.HashDay {
		m.InvHash = map[string]string{}
		for name, p := range o.state.Sent.fields() {
			if *p != nil {
				m.InvHash[name] = HashItems(*p)
			}
		}
		o.state.HashDay = day
		o.save()
	}
	return m
}

// Event returns the pending event message (tripped rules, inventory changes, web root
// changes), or nil when there is nothing or the last one went out less than gap ago.
func (o *Outbox) Event(now time.Time, gap time.Duration) *EventMsg {
	if now.Sub(o.lastEvt) < gap {
		return nil
	}
	c := o.c
	var diff map[string]SetDiff
	if o.haveFull {
		diff = DiffInv(o.state.Sent, o.cur)
	}
	c.mu.Lock()
	web := c.webPend
	c.mu.Unlock()
	if len(c.pending) == 0 && len(diff) == 0 && web == nil {
		return nil
	}
	m := &EventMsg{Header: o.header("event", now), Events: c.pending, Web: web}
	if len(diff) > 0 {
		m.InvDiff = diff
		if d, ok := diff["ports"]; ok {
			m.Ports = o.portDetails(d.Added)
		}
		for name, p := range o.cur.fields() {
			if *p != nil {
				*o.state.Sent.fields()[name] = *p
			}
		}
		o.save()
	}
	c.pending = nil
	c.mu.Lock()
	c.webPend = nil
	c.mu.Unlock()
	o.lastEvt = now
	return m
}

func (h Header) kind() string { return h.Type }
rules.go194 linesraw
package main

import (
	"sort"
	"time"
)

// Local rules decide on the server when something is worth reporting at once (an "event"
// message) instead of waiting for the 15-minute rollup. Each rule has a trip condition, a
// clear condition with hysteresis and a minimum duration, so a single odd sample never trips it.
const (
	diskBad, diskGood    = 90.0, 87.0
	fullSoonH, fullOkH   = 48.0, 72.0
	memBad, memGood      = 95.0, 90.0
	swapBad, swapGood    = 500.0, 100.0 // pages swapped in + out per second
	loadFactor, loadGood = 2.0, 1.5     // x CPU cores
	memFor, swapFor      = 10 * time.Minute, 10 * time.Minute
	loadFor              = 15 * time.Minute
	diskFor              = 2 * time.Minute // two samples in a row at the default interval
	clearFor             = 5 * time.Minute
	signalGoneFor        = 5 * time.Minute
)

// Event is one tripped rule, with its evidence.
type Event struct {
	Kind   string         `json:"kind"` // disk | inode | memory | swap | load | miner | tmp_exe | deleted_exe | cpu_hog
	Key    string         `json:"key"`
	At     int64          `json:"at"`
	Detail map[string]any `json:"detail"`
}

type ruleState struct {
	badSince  time.Time // zero: not bad now
	goodSince time.Time
	active    bool
}

// Rules keeps the state of every rule key ("disk:/", "memory", "miner:xmrig:/tmp/xmrig", ...).
type Rules struct {
	st       map[string]*ruleState
	lastSeen map[string]time.Time // process signals: last sample that showed them
}

func NewRules() *Rules { return &Rules{st: map[string]*ruleState{}, lastSeen: map[string]time.Time{}} }

// step advances one key. bad/good are this sample's readings; it returns true when the key
// trips now (it was not active) and false otherwise; cleared reports the reverse.
func (r *Rules) step(key string, now time.Time, bad, good bool, need time.Duration) (tripped, cleared bool) {
	s := r.st[key]
	if s == nil {
		s = &ruleState{}
		r.st[key] = s
	}
	if bad {
		if s.badSince.IsZero() {
			s.badSince = now
		}
	} else {
		s.badSince = time.Time{}
	}
	if good {
		if s.goodSince.IsZero() {
			s.goodSince = now
		}
	} else {
		s.goodSince = time.Time{}
	}
	if !s.active && bad && now.Sub(s.badSince) >= need-time.Second {
		s.active = true
		return true, false
	}
	if s.active && good && now.Sub(s.goodSince) >= clearFor-time.Second {
		s.active = false
		return false, true
	}
	return false, false
}

// Reading is what the rules look at in one sample.
type Reading struct {
	Mem, SwapIO, SwapPct, Load1 float64
	Cores                       int
	Disks                       []DiskOut
	TopCPU, TopMem              []Proc
	Signals                     map[string]Signal // process signals seen in this sample
}

// Update runs every rule on one sample. Returns the events that tripped now and the keys
// that cleared now.
func (r *Rules) Update(now time.Time, in Reading) (events []Event, cleared []string) {
	add := func(kind, key string, detail map[string]any) {
		events = append(events, Event{Kind: kind, Key: key, At: now.Unix(), Detail: detail})
	}
	mounts := map[string]bool{}
	for _, d := range in.Disks {
		mounts[d.Mount] = true
		soon := d.FullInH != nil && *d.FullInH < fullSoonH
		good := d.Pct < diskGood && (d.FullInH == nil || *d.FullInH >= fullOkH)
		if t, c := r.step("disk:"+d.Mount, now, d.Pct >= diskBad || soon, good, diskFor); t {
			add("disk", "disk:"+d.Mount, map[string]any{"mount": d.Mount, "pct": d.Pct, "full_in_h": d.FullInH, "rate_bph": d.RateBPH,
				"total": d.Total, "used": d.Used})
		} else if c {
			cleared = append(cleared, "disk:"+d.Mount)
		}
		if t, c := r.step("inode:"+d.Mount, now, d.InodesPct >= diskBad, d.InodesPct < diskGood, diskFor); t {
			add("inode", "inode:"+d.Mount, map[string]any{"mount": d.Mount, "pct": d.InodesPct})
		} else if c {
			cleared = append(cleared, "inode:"+d.Mount)
		}
	}
	for key, s := range r.st { // a filesystem that went away ends its rules
		for _, p := range []string{"disk:", "inode:"} {
			if len(key) > len(p) && key[:len(p)] == p && !mounts[key[len(p):]] {
				if s.active {
					cleared = append(cleared, key)
				}
				delete(r.st, key)
			}
		}
	}
	top := func(ps []Proc) any {
		if len(ps) == 0 {
			return nil
		}
		return ps[0]
	}
	minutes := func(key string) int { return int(now.Sub(r.st[key].badSince).Minutes() + 0.5) }
	if t, c := r.step("memory", now, in.Mem >= memBad, in.Mem < memGood, memFor); t {
		add("memory", "memory", map[string]any{"pct": in.Mem, "minutes": minutes("memory"), "top": top(in.TopMem)})
	} else if c {
		cleared = append(cleared, "memory")
	}
	if t, c := r.step("swap", now, in.SwapIO >= swapBad, in.SwapIO < swapGood, swapFor); t {
		add("swap", "swap", map[string]any{"rate": in.SwapIO, "pct": in.SwapPct, "minutes": minutes("swap")})
	} else if c {
		cleared = append(cleared, "swap")
	}
	cores := float64(maxInt(in.Cores, 1))
	if t, c := r.step("load", now, in.Load1 > loadFactor*cores, in.Load1 < loadGood*cores, loadFor); t {
		add("load", "load", map[string]any{"load": in.Load1, "cores": in.Cores, "limit": loadFactor * cores, "minutes": minutes("load"),
			"top": top(in.TopCPU)})
	} else if c {
		cleared = append(cleared, "load")
	}
	// Process signals: active while seen, cleared after 5 minutes without them.
	keys := make([]string, 0, len(in.Signals))
	for k := range in.Signals {
		keys = append(keys, k)
	}
	sort.Strings(keys)
	for _, k := range keys {
		sig := in.Signals[k]
		r.lastSeen[k] = now
		s := r.st[k]
		if s == nil || !s.active {
			r.st[k] = &ruleState{active: true, badSince: now}
			add(sig.Kind, k, sig.Detail)
		}
	}
	for k, seen := range r.lastSeen {
		if _, now2 := in.Signals[k]; now2 {
			continue
		}
		if now.Sub(seen) >= signalGoneFor {
			delete(r.lastSeen, k)
			if s := r.st[k]; s != nil && s.active {
				cleared = append(cleared, k)
			}
			delete(r.st, k)
		}
	}
	sort.Strings(cleared)
	return events, cleared
}

// Active lists the keys that are tripped now.
func (r *Rules) Active() []string {
	out := []string{}
	for k, s := range r.st {
		if s.active {
			out = append(out, k)
		}
	}
	sort.Strings(out)
	return out
}

func maxInt(a, b int) int {
	if a > b {
		return a
	}
	return b
}
security.go520 linesraw
package main

// Security signals. All of them are heuristics: a match is a reason to look, never proof,
// and the absence of a match proves nothing either. The agent only reads: it never kills,
// quarantines or changes anything.

import (
	"bufio"
	"bytes"
	"crypto/sha256"
	_ "embed"
	"encoding/base64"
	"errors"
	"io/fs"
	"os"
	"path/filepath"
	"sort"
	"strings"
)

//go:embed signatures.txt
var signaturesTxt []byte

// Signatures from signatures.txt.
type Signatures struct {
	Names map[string]bool
	Args  []string
	Busy  map[string]bool
}

func ParseSignatures(b []byte) Signatures {
	s := Signatures{Names: map[string]bool{}, Busy: map[string]bool{}}
	sc := bufio.NewScanner(bytes.NewReader(b))
	for sc.Scan() {
		line := strings.TrimSpace(sc.Text())
		if line == "" || strings.HasPrefix(line, "#") {
			continue
		}
		kind, val, ok := strings.Cut(line, " ")
		val = strings.ToLower(strings.TrimSpace(val))
		if !ok || val == "" {
			continue
		}
		switch kind {
		case "name":
			s.Names[val] = true
		case "arg":
			s.Args = append(s.Args, val)
		case "busy":
			s.Busy[val] = true
		}
	}
	return s
}

var sigs = ParseSignatures(signaturesTxt)

// Programs that show other programs' names in their arguments (someone looking for a
// miner with `grep xmrig` is not running one).
var lookers = map[string]bool{
	"grep": true, "egrep": true, "fgrep": true, "rg": true, "ag": true, "pgrep": true, "pkill": true, "ps": true,
	"less": true, "more": true, "tail": true, "head": true, "cat": true, "vi": true, "vim": true, "nano": true,
	"emacs": true, "journalctl": true, "man": true, "find": true, "locate": true, "htop": true, "top": true,
	"approvalens-age": true, "approvalens-agent": true,
}

// MinerMatch: why a process looks like a crypto miner ("" = it does not).
func (s Signatures) MinerMatch(comm, argv0, cmdline string) string {
	c := strings.ToLower(comm)
	a0 := strings.ToLower(filepath.Base(argv0))
	if s.Names[c] {
		return "name:" + c
	}
	if argv0 != "" && s.Names[a0] {
		return "name:" + a0
	}
	if lookers[c] || lookers[a0] {
		return ""
	}
	lc := strings.ToLower(cmdline)
	for _, a := range s.Args {
		if strings.Contains(lc, a) {
			return "arg:" + a
		}
	}
	return ""
}

var tmpDirs = []string{"/tmp/", "/var/tmp/", "/dev/shm/", "/run/shm/"}

// InTmp: an executable that runs from a world-writable temporary directory.
func InTmp(path string) bool {
	for _, d := range tmpDirs {
		if strings.HasPrefix(path, d) {
			return true
		}
	}
	return false
}

// SuspiciousDeleted: the process runs a binary that was deleted after it started. A
// package upgrade does that to every running daemon, so binaries under the system
// directories are left out; a memfd (a program that never touched the disk) is not.
func SuspiciousDeleted(exe string) bool {
	if !strings.HasSuffix(exe, " (deleted)") {
		return false
	}
	p := strings.TrimSuffix(exe, " (deleted)")
	if strings.HasPrefix(p, "/memfd:") {
		return true
	}
	for _, d := range []string{"/usr/", "/bin/", "/sbin/", "/lib/", "/lib64/", "/snap/", "/opt/", "/nix/", "/var/lib/docker/"} {
		if strings.HasPrefix(p, d) {
			return false
		}
	}
	return true
}

// ---------------------------------------------------------------------------
// crontabs
// ---------------------------------------------------------------------------

// ParseCrontab returns the active lines of one crontab ("<source>: <line>"), whitespace
// collapsed, comments and blank lines dropped. Variable lines (PATH=..., MAILTO=...) are
// kept: a changed PATH is a known way to hijack a job.
func ParseCrontab(content []byte, source string) []string {
	var out []string
	sc := bufio.NewScanner(bytes.NewReader(content))
	sc.Buffer(make([]byte, 64*1024), 1024*1024)
	for sc.Scan() {
		line := strings.Join(strings.Fields(sc.Text()), " ")
		if line == "" || strings.HasPrefix(line, "#") {
			continue
		}
		out = append(out, truncate(source+": "+line, 300))
	}
	return out
}

// DiffSets returns the items of cur that are not in prev, and the items of prev that are
// not in cur (both sorted).
func DiffSets(prev, cur []string) (added, removed []string) {
	p := make(map[string]bool, len(prev))
	for _, x := range prev {
		p[x] = true
	}
	c := make(map[string]bool, len(cur))
	for _, x := range cur {
		c[x] = true
		if !p[x] {
			added = append(added, x)
		}
	}
	for _, x := range prev {
		if !c[x] {
			removed = append(removed, x)
		}
	}
	sort.Strings(added)
	sort.Strings(removed)
	return dedupSorted(added), dedupSorted(removed)
}

func dedupSorted(xs []string) []string {
	if len(xs) < 2 {
		return xs
	}
	out := xs[:1]
	for _, x := range xs[1:] {
		if x != out[len(out)-1] {
			out = append(out, x)
		}
	}
	return out
}

// ReadCrontabs collects every system and user crontab it can read. unreadable lists the
// places it could not (other users' crontabs need root).
func ReadCrontabs(root string) (entries []string, unreadable []string) {
	files := []string{filepath.Join(root, "etc/crontab")}
	for _, dir := range []string{"etc/cron.d", "var/spool/cron/crontabs", "var/spool/cron"} {
		d := filepath.Join(root, dir)
		ents, err := os.ReadDir(d)
		if err != nil {
			if errors.Is(err, fs.ErrPermission) {
				unreadable = append(unreadable, "/"+dir)
			}
			continue
		}
		for _, e := range ents {
			if e.Type().IsRegular() && !strings.HasPrefix(e.Name(), ".") {
				files = append(files, filepath.Join(d, e.Name()))
			}
		}
	}
	for _, f := range files {
		b, err := readSmall(f, 256*1024)
		if err != nil {
			if errors.Is(err, fs.ErrPermission) {
				unreadable = append(unreadable, strings.TrimPrefix(f, strings.TrimSuffix(root, "/")))
			}
			continue
		}
		entries = append(entries, ParseCrontab(b, strings.TrimPrefix(f, strings.TrimSuffix(root, "/")))...)
	}
	sort.Strings(entries)
	entries = dedupSorted(entries)
	if len(entries) > 300 {
		entries = entries[:300]
	}
	return entries, unreadable
}

// ---------------------------------------------------------------------------
// SSH authorized keys
// ---------------------------------------------------------------------------

// ParseAuthorizedKeys returns "<user> <type> SHA256:<fingerprint> <comment>" per key.
func ParseAuthorizedKeys(b []byte, user string) []string {
	var out []string
	sc := bufio.NewScanner(bytes.NewReader(b))
	sc.Buffer(make([]byte, 64*1024), 1024*1024)
	for sc.Scan() {
		line := strings.TrimSpace(sc.Text())
		if line == "" || strings.HasPrefix(line, "#") {
			continue
		}
		f := strings.Fields(line)
		idx := -1
		for i, x := range f {
			if keyType(x) {
				idx = i
				break
			}
		}
		if idx < 0 || idx+1 >= len(f) {
			continue
		}
		blob, err := base64.StdEncoding.DecodeString(f[idx+1])
		if err != nil {
			continue
		}
		sum := sha256.Sum256(blob)
		item := user + " " + f[idx] + " SHA256:" + base64.RawStdEncoding.EncodeToString(sum[:])
		if idx > 0 {
			item += " [options]"
		}
		if idx+2 < len(f) {
			item += " " + truncate(strings.Join(f[idx+2:], " "), 60)
		}
		out = append(out, item)
	}
	return out
}

func keyType(s string) bool {
	return strings.HasPrefix(s, "ssh-") || strings.HasPrefix(s, "ecdsa-sha2-") || strings.HasPrefix(s, "sk-ssh-") ||
		strings.HasPrefix(s, "sk-ecdsa-")
}

// ReadSSHKeys reads authorized_keys of root and of every account with a login shell.
func ReadSSHKeys(root string, users []PasswdEntry) (keys []string, unreadable []string) {
	seen := map[string]bool{}
	for _, u := range users {
		if u.UID != 0 && !LoginShell(u.Shell) {
			continue
		}
		if u.Home == "" || u.Home == "/" || seen[u.Home] {
			continue
		}
		seen[u.Home] = true
		for _, name := range []string{"authorized_keys", "authorized_keys2"} {
			p := filepath.Join(root, u.Home, ".ssh", name)
			b, err := readSmall(p, 512*1024)
			if err != nil {
				if errors.Is(err, fs.ErrPermission) {
					unreadable = append(unreadable, filepath.Join(u.Home, ".ssh", name))
				}
				continue
			}
			keys = append(keys, ParseAuthorizedKeys(b, u.Name)...)
		}
	}
	sort.Strings(keys)
	keys = dedupSorted(keys)
	if len(keys) > 200 {
		keys = keys[:200]
	}
	return keys, dedupSorted(sortedCopy(unreadable))
}

// ---------------------------------------------------------------------------
// setuid / setgid binaries
// ---------------------------------------------------------------------------

var setuidRoots = []string{"/bin", "/sbin", "/usr/bin", "/usr/sbin", "/usr/local/bin", "/usr/local/sbin", "/usr/lib",
	"/usr/libexec", "/usr/local/lib", "/opt", "/tmp", "/var/tmp", "/dev/shm"}

// ScanSetuid walks the usual binary directories (and the temporary ones) for setuid or
// setgid files: "<path> <mode> uid=<n>". It stays light: no symlinks are followed and
// it stops after maxFiles entries.
func ScanSetuid(roots []string, maxFiles int) (items []string, truncated bool) {
	seen := map[string]bool{}
	count := 0
	for _, r := range roots {
		real, err := filepath.EvalSymlinks(r)
		if err != nil || seen[real] {
			continue
		}
		seen[real] = true
		_ = filepath.WalkDir(real, func(p string, d fs.DirEntry, err error) error {
			if err != nil {
				if d != nil && d.IsDir() {
					return fs.SkipDir
				}
				return nil
			}
			count++
			if count > maxFiles {
				truncated = true
				return fs.SkipAll
			}
			if d.IsDir() {
				if strings.Count(strings.TrimPrefix(p, real), "/") > 8 {
					return fs.SkipDir
				}
				return nil
			}
			if !d.Type().IsRegular() {
				return nil
			}
			info, err := d.Info()
			if err != nil {
				return nil
			}
			m := info.Mode()
			if m&(fs.ModeSetuid|fs.ModeSetgid) == 0 {
				return nil
			}
			items = append(items, p+" "+modeString(m)+" uid="+itoa(fileUID(info)))
			return nil
		})
		if truncated {
			break
		}
	}
	sort.Strings(items)
	return items, truncated
}

func modeString(m fs.FileMode) string {
	v := uint32(m.Perm())
	if m&fs.ModeSetuid != 0 {
		v |= 0o4000
	}
	if m&fs.ModeSetgid != 0 {
		v |= 0o2000
	}
	if m&fs.ModeSticky != 0 {
		v |= 0o1000
	}
	s := []byte("0000")
	for i := 3; i >= 0; i-- {
		s[i] = byte('0' + v&7)
		v >>= 3
	}
	return string(s)
}

// ---------------------------------------------------------------------------
// web roots
// ---------------------------------------------------------------------------

var webExts = map[string]bool{".php": true, ".phtml": true, ".php3": true, ".php4": true, ".php5": true, ".php7": true,
	".phar": true, ".inc": true, ".js": true, ".mjs": true, ".html": true, ".htm": true, ".shtml": true, ".htaccess": true,
	".user.ini": true}

var webSkipDirs = map[string]bool{"node_modules": true, ".git": true, ".svn": true, ".hg": true, "cache": true, ".cache": true}

// WebFile key: modification time and size, never the content.
func webWatched(name string) bool {
	if name == ".htaccess" || name == ".user.ini" {
		return true
	}
	return webExts[strings.ToLower(filepath.Ext(name))]
}

// ScanWebRoots maps each PHP / JS / HTML file (and .htaccess) under the roots to
// "<mtime>:<size>". Stops at maxFiles (truncated=true).
func ScanWebRoots(roots []string, exclude []string, maxFiles int) (files map[string]string, truncated bool) {
	files = map[string]string{}
	skip := map[string]bool{}
	for k := range webSkipDirs {
		skip[k] = true
	}
	for _, x := range exclude {
		skip[x] = true
	}
	for _, r := range roots {
		_ = filepath.WalkDir(r, func(p string, d fs.DirEntry, err error) error {
			if err != nil {
				if d != nil && d.IsDir() {
					return fs.SkipDir
				}
				return nil
			}
			if d.IsDir() {
				if p != r && (skip[d.Name()] || skip[p]) {
					return fs.SkipDir
				}
				return nil
			}
			if !d.Type().IsRegular() || !webWatched(d.Name()) {
				return nil
			}
			if len(files) >= maxFiles {
				truncated = true
				return fs.SkipAll
			}
			info, err := d.Info()
			if err != nil {
				return nil
			}
			files[p] = itoa64(info.ModTime().Unix()) + ":" + itoa64(info.Size())
			return nil
		})
		if truncated {
			break
		}
	}
	return files, truncated
}

// WebChanges between two scans: new, changed and removed file names (capped).
type WebChanges struct {
	Added     []string `json:"added"`
	Changed   []string `json:"changed"`
	Removed   []string `json:"removed"`
	NAdded    int      `json:"n_added"`
	NChanged  int      `json:"n_changed"`
	NRemoved  int      `json:"n_removed"`
	Files     int      `json:"files"`
	Truncated bool     `json:"truncated,omitempty"`
}

func (w WebChanges) Empty() bool { return w.NAdded+w.NChanged+w.NRemoved == 0 }

func DiffWebRoots(prev, cur map[string]string, limit int) WebChanges {
	var w WebChanges
	w.Files = len(cur)
	for p, v := range cur {
		old, ok := prev[p]
		switch {
		case !ok:
			w.NAdded++
			w.Added = append(w.Added, p)
		case old != v:
			w.NChanged++
			w.Changed = append(w.Changed, p)
		}
	}
	for p := range prev {
		if _, ok := cur[p]; !ok {
			w.NRemoved++
			w.Removed = append(w.Removed, p)
		}
	}
	for _, xs := range []*[]string{&w.Added, &w.Changed, &w.Removed} {
		sort.Strings(*xs)
		if len(*xs) > limit {
			*xs = (*xs)[:limit]
		}
	}
	return w
}

// ---------------------------------------------------------------------------
// helpers
// ---------------------------------------------------------------------------

func readSmall(path string, max int64) ([]byte, error) {
	f, err := os.Open(path)
	if err != nil {
		return nil, err
	}
	defer f.Close()
	buf := make([]byte, 0, 4096)
	tmp := make([]byte, 32*1024)
	for int64(len(buf)) < max {
		n, err := f.Read(tmp)
		buf = append(buf, tmp[:n]...)
		if err != nil {
			break
		}
	}
	if int64(len(buf)) > max {
		buf = buf[:max]
	}
	return buf, nil
}

func truncate(s string, n int) string {
	if len(s) <= n {
		return s
	}
	// keep valid UTF-8: back off to a rune start
	cut := n
	for cut > 0 && (s[cut]&0xC0) == 0x80 {
		cut--
	}
	return s[:cut] + "…"
}

func sortedCopy(xs []string) []string {
	out := append([]string(nil), xs...)
	sort.Strings(out)
	return out
}
procs.go248 linesraw
package main

import (
	"bytes"
	"os"
	"path/filepath"
	"regexp"
	"sort"
	"strconv"
	"strings"
	"sync"
	"syscall"
	"time"
)

const clockTicks = 100 // USER_HZ: 100 on every Linux build for amd64 and arm64

var pageSize = int64(os.Getpagesize())

// Proc is one process as reported (top lists and signal evidence).
type Proc struct {
	PID  int     `json:"pid"`
	Name string  `json:"name"`
	User string  `json:"user"`
	Exe  string  `json:"exe,omitempty"`
	Cmd  string  `json:"cmd,omitempty"`
	CPU  float64 `json:"cpu"` // percent of one core over the last sample interval
	RSS  int64   `json:"rss"` // bytes
	// internal
	uid       int
	argv0     string
	exeKnown  bool
	startTime uint64
	kernel    bool
}

type procPrev struct {
	ticks uint64
	start uint64
}

// ProcSampler reads every /proc/<pid> once per sample and keeps the CPU ticks of the
// previous sample to turn them into a percentage.
type ProcSampler struct {
	root  string
	prev  map[int]procPrev
	prevT time.Time
	users *UserCache
}

func NewProcSampler(root string, users *UserCache) *ProcSampler {
	return &ProcSampler{root: root, prev: map[int]procPrev{}, users: users}
}

func (s *ProcSampler) Sample(now time.Time) []Proc {
	ents, err := os.ReadDir(s.root)
	if err != nil {
		return nil
	}
	elapsed := now.Sub(s.prevT).Seconds()
	next := make(map[int]procPrev, len(s.prev))
	out := make([]Proc, 0, len(ents))
	for _, e := range ents {
		pid, err := strconv.Atoi(e.Name())
		if err != nil || pid <= 0 {
			continue
		}
		dir := filepath.Join(s.root, e.Name())
		b, err := os.ReadFile(filepath.Join(dir, "stat"))
		if err != nil {
			continue
		}
		st, err := ParsePidStat(b)
		if err != nil {
			continue
		}
		p := Proc{PID: pid, Name: st.Comm, RSS: st.RSSPages * pageSize, startTime: st.StartTime}
		p.kernel = pid == 2 || st.PPid == 2
		if fi, err := os.Stat(dir); err == nil {
			p.uid = fileUID(fi)
		}
		p.User = s.users.Name(p.uid)
		ticks := st.UTime + st.STime
		if pp, ok := s.prev[pid]; ok && pp.start == st.StartTime && elapsed > 0 && ticks >= pp.ticks {
			p.CPU = round1(100 * float64(ticks-pp.ticks) / clockTicks / elapsed)
		}
		next[pid] = procPrev{ticks: ticks, start: st.StartTime}
		if !p.kernel {
			if exe, err := os.Readlink(filepath.Join(dir, "exe")); err == nil {
				p.Exe, p.exeKnown = exe, true
			}
			if cb, err := readSmall(filepath.Join(dir, "cmdline"), 4096); err == nil && len(cb) > 0 {
				args := bytes.Split(bytes.TrimRight(cb, "\x00"), []byte{0})
				if len(args) > 0 {
					p.argv0 = string(args[0])
				}
				p.Cmd = string(bytes.Join(args, []byte{' '}))
			}
		}
		out = append(out, p)
	}
	s.prev = next
	s.prevT = now
	return out
}

// Reported strips what the report must not carry raw: the command line is redacted
// and cut to 200 characters.
func (p Proc) Reported() Proc {
	q := p
	q.Cmd = truncate(Redact(p.Cmd), 200)
	q.Exe = truncate(p.Exe, 200)
	return q
}

var (
	redactKV   = regexp.MustCompile(`(?i)((?:pass(?:word|wd)?|pwd|secret|token|api[_-]?key|apikey|access[_-]?key|auth|credentials?|private[_-]?key)[=:])\S+`)
	redactFlag = regexp.MustCompile(`(?i)(--(?:pass(?:word|wd)?|secret|token|api-?key|auth)[ =])\S+`)
	redactURL  = regexp.MustCompile(`(://[^/\s:@]+:)[^@\s/]+@`)
	redactMy   = regexp.MustCompile(`(\s-p)[^\s]+`)
)

// Redact hides values that look like secrets in a command line: key=value pairs named
// like passwords or tokens, --password values, user:password@ in URLs and mysql-style -pSECRET.
func Redact(cmd string) string {
	cmd = redactKV.ReplaceAllString(cmd, "${1}***")
	cmd = redactFlag.ReplaceAllString(cmd, "${1}***")
	cmd = redactURL.ReplaceAllString(cmd, "${1}***@")
	if strings.Contains(cmd, "mysql") || strings.Contains(cmd, "mariadb") {
		cmd = redactMy.ReplaceAllString(cmd, "${1}***")
	}
	return cmd
}

// TopBy returns the n processes with the highest key (ties: lower pid first).
func TopBy(ps []Proc, n int, key func(Proc) float64) []Proc {
	cp := append([]Proc(nil), ps...)
	sort.SliceStable(cp, func(i, j int) bool {
		ki, kj := key(cp[i]), key(cp[j])
		if ki != kj {
			return ki > kj
		}
		return cp[i].PID < cp[j].PID
	})
	if len(cp) > n {
		cp = cp[:n]
	}
	out := make([]Proc, len(cp))
	for i, p := range cp {
		out[i] = p.Reported()
	}
	return out
}

// UserCache maps uids to names from /etc/passwd, reloading when the file changes.
type UserCache struct {
	path  string
	mu    sync.Mutex
	mtime time.Time
	names map[int]string
	all   []PasswdEntry
}

func NewUserCache(path string) *UserCache { return &UserCache{path: path, names: map[int]string{}} }

func (u *UserCache) refresh() {
	fi, err := os.Stat(u.path)
	if err != nil || fi.ModTime().Equal(u.mtime) {
		return
	}
	b, err := os.ReadFile(u.path)
	if err != nil {
		return
	}
	u.mtime = fi.ModTime()
	u.all = ParsePasswd(b)
	u.names = make(map[int]string, len(u.all))
	for _, e := range u.all {
		if _, ok := u.names[e.UID]; !ok {
			u.names[e.UID] = e.Name
		}
	}
}

func (u *UserCache) Name(uid int) string {
	u.mu.Lock()
	defer u.mu.Unlock()
	u.refresh()
	if n, ok := u.names[uid]; ok {
		return n
	}
	return strconv.Itoa(uid)
}

func (u *UserCache) All() []PasswdEntry {
	u.mu.Lock()
	defer u.mu.Unlock()
	u.refresh()
	return append([]PasswdEntry(nil), u.all...)
}

func fileUID(fi os.FileInfo) int {
	if st, ok := fi.Sys().(*syscall.Stat_t); ok {
		return int(st.Uid)
	}
	return -1
}

// SocketOwners maps socket inodes to the process holding them, by reading the fd links
// of every process it may read (all of them as root, only its own otherwise).
func SocketOwners(root string, inodes map[uint64]bool, procs []Proc) map[uint64]Proc {
	out := map[uint64]Proc{}
	if len(inodes) == 0 {
		return out
	}
	for _, p := range procs {
		if p.kernel {
			continue
		}
		fdDir := filepath.Join(root, strconv.Itoa(p.PID), "fd")
		ents, err := os.ReadDir(fdDir)
		if err != nil {
			continue
		}
		for _, e := range ents {
			l, err := os.Readlink(filepath.Join(fdDir, e.Name()))
			if err != nil || !strings.HasPrefix(l, "socket:[") {
				continue
			}
			ino, err := strconv.ParseUint(strings.TrimSuffix(strings.TrimPrefix(l, "socket:["), "]"), 10, 64)
			if err == nil && inodes[ino] {
				if _, seen := out[ino]; !seen {
					out[ino] = p
				}
			}
		}
		if len(out) == len(inodes) {
			break
		}
	}
	return out
}

func round1(v float64) float64 { return float64(int64(v*10+0.5)) / 10 }

func itoa(n int) string     { return strconv.Itoa(n) }
func itoa64(n int64) string { return strconv.FormatInt(n, 10) }
procfs.go428 linesraw
package main

// Parsers for the plain-text files under /proc. Each parser takes the file's
// bytes so the tests can feed fixtures (testdata/); the readers that open the
// real files live next to the code that samples them.

import (
	"bufio"
	"bytes"
	"encoding/hex"
	"errors"
	"fmt"
	"net"
	"strconv"
	"strings"
)

// CPUTimes is the aggregate "cpu" line of /proc/stat, in clock ticks.
type CPUTimes struct {
	User, Nice, System, Idle, IOWait, IRQ, SoftIRQ, Steal uint64
}

func (c CPUTimes) Total() uint64 {
	return c.User + c.Nice + c.System + c.Idle + c.IOWait + c.IRQ + c.SoftIRQ + c.Steal
}

// Busy is everything but idle and iowait.
func (c CPUTimes) Busy() uint64 { return c.Total() - c.Idle - c.IOWait }

// ProcStat is what the sampler needs from /proc/stat.
type ProcStat struct {
	CPU      CPUTimes
	Cores    int
	BootTime int64
}

func ParseProcStat(b []byte) (ProcStat, error) {
	var ps ProcStat
	found := false
	sc := bufio.NewScanner(bytes.NewReader(b))
	sc.Buffer(make([]byte, 64*1024), 1024*1024)
	for sc.Scan() {
		f := strings.Fields(sc.Text())
		if len(f) == 0 {
			continue
		}
		switch {
		case f[0] == "cpu":
			if len(f) < 5 {
				return ps, errors.New("short cpu line")
			}
			v := make([]uint64, 8)
			for i := 0; i < 8 && i+1 < len(f); i++ {
				n, err := strconv.ParseUint(f[i+1], 10, 64)
				if err != nil {
					return ps, fmt.Errorf("cpu field %d: %w", i, err)
				}
				v[i] = n
			}
			ps.CPU = CPUTimes{v[0], v[1], v[2], v[3], v[4], v[5], v[6], v[7]}
			found = true
		case strings.HasPrefix(f[0], "cpu"):
			ps.Cores++
		case f[0] == "btime" && len(f) > 1:
			ps.BootTime, _ = strconv.ParseInt(f[1], 10, 64)
		}
	}
	if !found {
		return ps, errors.New("no cpu line")
	}
	if ps.Cores == 0 {
		ps.Cores = 1
	}
	return ps, nil
}

// CPUPercent between two readings of the aggregate line (0..100).
func CPUPercent(prev, cur CPUTimes) float64 {
	dt := float64(cur.Total()) - float64(prev.Total())
	if dt <= 0 {
		return 0
	}
	db := float64(cur.Busy()) - float64(prev.Busy())
	if db < 0 {
		db = 0
	}
	return clamp(100*db/dt, 0, 100)
}

// MemInfo in bytes.
type MemInfo struct {
	Total, Available, Free, Buffers, Cached, SwapTotal, SwapFree uint64
}

func (m MemInfo) UsedPct() float64 {
	if m.Total == 0 {
		return 0
	}
	return clamp(100*float64(m.Total-minU(m.Available, m.Total))/float64(m.Total), 0, 100)
}

func (m MemInfo) SwapUsed() uint64 {
	if m.SwapFree > m.SwapTotal {
		return 0
	}
	return m.SwapTotal - m.SwapFree
}

func (m MemInfo) SwapPct() float64 {
	if m.SwapTotal == 0 {
		return 0
	}
	return clamp(100*float64(m.SwapUsed())/float64(m.SwapTotal), 0, 100)
}

func ParseMemInfo(b []byte) (MemInfo, error) {
	var m MemInfo
	hasAvail := false
	sc := bufio.NewScanner(bytes.NewReader(b))
	for sc.Scan() {
		line := sc.Text()
		i := strings.IndexByte(line, ':')
		if i < 0 {
			continue
		}
		f := strings.Fields(line[i+1:])
		if len(f) == 0 {
			continue
		}
		n, err := strconv.ParseUint(f[0], 10, 64)
		if err != nil {
			continue
		}
		if len(f) > 1 && f[1] == "kB" {
			n *= 1024
		}
		switch line[:i] {
		case "MemTotal":
			m.Total = n
		case "MemAvailable":
			m.Available, hasAvail = n, true
		case "MemFree":
			m.Free = n
		case "Buffers":
			m.Buffers = n
		case "Cached":
			m.Cached = n
		case "SwapTotal":
			m.SwapTotal = n
		case "SwapFree":
			m.SwapFree = n
		}
	}
	if m.Total == 0 {
		return m, errors.New("no MemTotal")
	}
	if !hasAvail { // kernels before 3.14
		m.Available = m.Free + m.Buffers + m.Cached
	}
	return m, nil
}

// ParseVMStat returns the swap-in / swap-out page counters.
func ParseVMStat(b []byte) (pswpin, pswpout uint64) {
	sc := bufio.NewScanner(bytes.NewReader(b))
	for sc.Scan() {
		f := strings.Fields(sc.Text())
		if len(f) != 2 {
			continue
		}
		switch f[0] {
		case "pswpin":
			pswpin, _ = strconv.ParseUint(f[1], 10, 64)
		case "pswpout":
			pswpout, _ = strconv.ParseUint(f[1], 10, 64)
		}
	}
	return
}

func ParseLoadAvg(b []byte) ([3]float64, error) {
	var l [3]float64
	f := strings.Fields(string(b))
	if len(f) < 3 {
		return l, errors.New("short loadavg")
	}
	for i := 0; i < 3; i++ {
		v, err := strconv.ParseFloat(f[i], 64)
		if err != nil {
			return l, err
		}
		l[i] = v
	}
	return l, nil
}

func ParseUptime(b []byte) (float64, error) {
	f := strings.Fields(string(b))
	if len(f) < 1 {
		return 0, errors.New("empty uptime")
	}
	return strconv.ParseFloat(f[0], 64)
}

// NetCounters are the summed byte counters of the interfaces that carry real traffic.
type NetCounters struct{ RX, TX uint64 }

// virtualIface: loopback and the host side of containers / bridges (their traffic is
// counted again on the physical interface).
func virtualIface(name string) bool {
	if name == "lo" {
		return true
	}
	for _, p := range []string{"veth", "docker", "br-", "virbr", "cni", "flannel", "cali", "vxlan", "tun", "tap", "kube", "lxc", "vnet"} {
		if strings.HasPrefix(name, p) {
			return true
		}
	}
	return false
}

func ParseNetDev(b []byte) NetCounters {
	var n NetCounters
	sc := bufio.NewScanner(bytes.NewReader(b))
	for sc.Scan() {
		line := sc.Text()
		i := strings.IndexByte(line, ':')
		if i < 0 {
			continue
		}
		name := strings.TrimSpace(line[:i])
		if virtualIface(name) {
			continue
		}
		f := strings.Fields(line[i+1:])
		if len(f) < 9 {
			continue
		}
		rx, _ := strconv.ParseUint(f[0], 10, 64)
		tx, _ := strconv.ParseUint(f[8], 10, 64)
		n.RX += rx
		n.TX += tx
	}
	return n
}

// Socket is one listening socket from /proc/net/{tcp,tcp6,udp,udp6}.
type Socket struct {
	Proto string // tcp | udp (the address family is in Addr)
	Addr  string
	Port  int
	UID   int
	Inode uint64
}

// Local: bound to a loopback address, so not reachable from outside.
func (s Socket) Local() bool {
	ip := net.ParseIP(s.Addr)
	return ip != nil && ip.IsLoopback()
}

// ParseNetSockets reads one /proc/net/{tcp,udp}[6] table and returns the listening
// sockets: TCP in state LISTEN (0A), UDP unconnected (07) with no remote peer.
func ParseNetSockets(b []byte, proto string) ([]Socket, error) {
	var out []Socket
	sc := bufio.NewScanner(bytes.NewReader(b))
	first := true
	for sc.Scan() {
		if first { // header
			first = false
			continue
		}
		f := strings.Fields(sc.Text())
		if len(f) < 10 {
			continue
		}
		st := f[3]
		if proto == "tcp" && st != "0A" {
			continue
		}
		if proto == "udp" {
			if st != "07" {
				continue
			}
			if _, rport, err := splitHexAddr(f[2]); err != nil || rport != 0 {
				continue
			}
		}
		addr, port, err := splitHexAddr(f[1])
		if err != nil {
			return nil, err
		}
		uid, _ := strconv.Atoi(f[7])
		inode, _ := strconv.ParseUint(f[9], 10, 64)
		out = append(out, Socket{Proto: proto, Addr: addr, Port: port, UID: uid, Inode: inode})
	}
	return out, nil
}

// splitHexAddr decodes "0100007F:1F90" (IPv4) or the 32-hex-digit IPv6 form. The kernel
// prints each 32-bit word in host byte order (little endian on amd64 and arm64).
func splitHexAddr(s string) (string, int, error) {
	i := strings.IndexByte(s, ':')
	if i < 0 {
		return "", 0, fmt.Errorf("bad address %q", s)
	}
	raw, err := hex.DecodeString(s[:i])
	if err != nil || (len(raw) != 4 && len(raw) != 16) {
		return "", 0, fmt.Errorf("bad address %q", s)
	}
	port, err := strconv.ParseUint(s[i+1:], 16, 16)
	if err != nil {
		return "", 0, fmt.Errorf("bad port %q", s)
	}
	ip := make(net.IP, len(raw))
	for w := 0; w < len(raw); w += 4 {
		ip[w], ip[w+1], ip[w+2], ip[w+3] = raw[w+3], raw[w+2], raw[w+1], raw[w]
	}
	return ip.String(), int(port), nil
}

// PidStat is what we use from /proc/<pid>/stat.
type PidStat struct {
	Comm      string
	State     byte
	PPid      int
	UTime     uint64
	STime     uint64
	StartTime uint64 // clock ticks after boot
	RSSPages  int64
}

// ParsePidStat handles a comm that contains spaces or parentheses: it runs from the
// first '(' to the last ')'.
func ParsePidStat(b []byte) (PidStat, error) {
	var p PidStat
	s := string(b)
	l, r := strings.IndexByte(s, '('), strings.LastIndexByte(s, ')')
	if l < 0 || r < l {
		return p, errors.New("bad stat")
	}
	p.Comm = s[l+1 : r]
	f := strings.Fields(s[r+1:])
	// f[0]=state (field 3) ... utime is field 14 -> f[11], stime f[12], starttime field 22 -> f[19], rss field 24 -> f[21]
	if len(f) < 22 {
		return p, errors.New("short stat")
	}
	p.State = f[0][0]
	p.PPid, _ = strconv.Atoi(f[1])
	p.UTime, _ = strconv.ParseUint(f[11], 10, 64)
	p.STime, _ = strconv.ParseUint(f[12], 10, 64)
	p.StartTime, _ = strconv.ParseUint(f[19], 10, 64)
	p.RSSPages, _ = strconv.ParseInt(f[21], 10, 64)
	return p, nil
}

// ParseOSRelease returns PRETTY_NAME (or NAME VERSION_ID).
func ParseOSRelease(b []byte) string {
	vals := map[string]string{}
	sc := bufio.NewScanner(bytes.NewReader(b))
	for sc.Scan() {
		k, v, ok := strings.Cut(strings.TrimSpace(sc.Text()), "=")
		if !ok {
			continue
		}
		vals[k] = strings.Trim(v, `"'`)
	}
	if v := vals["PRETTY_NAME"]; v != "" {
		return v
	}
	return strings.TrimSpace(vals["NAME"] + " " + vals["VERSION_ID"])
}

// Users from /etc/passwd: uid -> name, plus the accounts that can log in.
type PasswdEntry struct {
	Name  string
	UID   int
	Home  string
	Shell string
}

func ParsePasswd(b []byte) []PasswdEntry {
	var out []PasswdEntry
	sc := bufio.NewScanner(bytes.NewReader(b))
	for sc.Scan() {
		f := strings.Split(sc.Text(), ":")
		if len(f) < 7 || strings.HasPrefix(f[0], "#") {
			continue
		}
		uid, err := strconv.Atoi(f[2])
		if err != nil {
			continue
		}
		out = append(out, PasswdEntry{Name: f[0], UID: uid, Home: f[5], Shell: f[6]})
	}
	return out
}

// LoginShell: an account whose shell lets someone log in.
func LoginShell(shell string) bool {
	if shell == "" {
		return false
	}
	base := shell[strings.LastIndexByte(shell, '/')+1:]
	switch base {
	case "nologin", "false", "sync", "halt", "shutdown", "true":
		return false
	}
	return true
}

func clamp(v, lo, hi float64) float64 {
	if v < lo {
		return lo
	}
	if v > hi {
		return hi
	}
	return v
}

func minU(a, b uint64) uint64 {
	if a < b {
		return a
	}
	return b
}
disk.go182 linesraw
package main

import (
	"bufio"
	"bytes"
	"sort"
	"strconv"
	"strings"
	"syscall"
)

// Mount is one real filesystem from /proc/self/mountinfo.
type Mount struct {
	Point  string
	FSType string
	Source string
}

// Filesystems that hold data. Pseudo filesystems, tmpfs, overlays, snap images
// (squashfs, always 100% full) and network filesystems (statfs can hang) are left out.
var realFS = map[string]bool{
	"ext2": true, "ext3": true, "ext4": true, "xfs": true, "btrfs": true, "zfs": true, "f2fs": true,
	"jfs": true, "reiserfs": true, "vfat": true, "exfat": true, "ntfs": true, "ntfs3": true,
	"bcachefs": true, "hfsplus": true, "ufs": true,
}

// unescapeMount decodes the octal escapes mountinfo uses (\040 = space).
func unescapeMount(s string) string {
	if !strings.Contains(s, `\`) {
		return s
	}
	var b strings.Builder
	for i := 0; i < len(s); i++ {
		if s[i] == '\\' && i+3 < len(s) {
			if n, err := strconv.ParseUint(s[i+1:i+4], 8, 8); err == nil {
				b.WriteByte(byte(n))
				i += 3
				continue
			}
		}
		b.WriteByte(s[i])
	}
	return b.String()
}

// ParseMountInfo returns one mount per real filesystem: a device mounted at several
// places (bind mounts, btrfs subvolumes, the read-only view systemd gives the agent)
// is listed once, at its shortest mount point.
func ParseMountInfo(b []byte) []Mount {
	best := map[string]Mount{}
	sc := bufio.NewScanner(bytes.NewReader(b))
	sc.Buffer(make([]byte, 64*1024), 1024*1024)
	for sc.Scan() {
		line := sc.Text()
		pre, post, ok := strings.Cut(line, " - ")
		if !ok {
			continue
		}
		pf, qf := strings.Fields(pre), strings.Fields(post)
		if len(pf) < 5 || len(qf) < 2 {
			continue
		}
		fstype, source := qf[0], qf[1]
		if !realFS[fstype] {
			continue
		}
		point := unescapeMount(pf[4])
		if skipMountPoint(point) {
			continue
		}
		key := source
		if !strings.HasPrefix(source, "/dev/") { // zfs datasets, odd sources: the device numbers
			key = pf[2] + "|" + source
		}
		cur, seen := best[key]
		if !seen || len(point) < len(cur.Point) {
			best[key] = Mount{Point: point, FSType: fstype, Source: source}
		}
	}
	out := make([]Mount, 0, len(best))
	for _, m := range best {
		out = append(out, m)
	}
	sort.Slice(out, func(i, j int) bool { return out[i].Point < out[j].Point })
	return out
}

func skipMountPoint(p string) bool {
	for _, pre := range []string{"/proc", "/sys", "/dev", "/run", "/snap", "/var/lib/docker", "/var/lib/containers", "/var/snap"} {
		if p == pre || strings.HasPrefix(p, pre+"/") {
			return true
		}
	}
	return false
}

// DiskUsage of one mount, in bytes and inodes.
type DiskUsage struct {
	Total, Used, Avail uint64
	Inodes, InodesFree uint64
	UsedPct, InodesPct float64
}

func statDisk(path string) (DiskUsage, error) {
	var st syscall.Statfs_t
	if err := syscall.Statfs(path, &st); err != nil {
		return DiskUsage{}, err
	}
	bs := uint64(st.Bsize)
	d := DiskUsage{Total: st.Blocks * bs, Avail: st.Bavail * bs, Inodes: st.Files, InodesFree: st.Ffree}
	free := st.Bfree * bs
	if d.Total >= free {
		d.Used = d.Total - free
	}
	// As df: used / (used + available to unprivileged users), so root's reserve counts as full.
	if den := d.Used + d.Avail; den > 0 {
		d.UsedPct = 100 * float64(d.Used) / float64(den)
	}
	if d.Inodes > 0 && d.Inodes >= d.InodesFree {
		d.InodesPct = 100 * float64(d.Inodes-d.InodesFree) / float64(d.Inodes)
	}
	return d, nil
}

// FillTracker keeps the used bytes of one mount over the last window (6 h) and projects
// when the disk is full if the trend goes on.
type FillTracker struct {
	Window int64 // seconds
	ts     []int64
	used   []float64
}

func NewFillTracker(window int64) *FillTracker { return &FillTracker{Window: window} }

func (f *FillTracker) Add(t int64, used uint64) {
	f.ts = append(f.ts, t)
	f.used = append(f.used, float64(used))
	cut := 0
	for cut < len(f.ts) && f.ts[cut] < t-f.Window {
		cut++
	}
	if cut > 0 {
		f.ts = append(f.ts[:0], f.ts[cut:]...)
		f.used = append(f.used[:0], f.used[cut:]...)
	}
}

// Rate is the least-squares slope in bytes per hour (0 without enough data: at least
// 30 minutes and 10 samples).
func (f *FillTracker) Rate() (float64, bool) {
	n := len(f.ts)
	if n < 10 || f.ts[n-1]-f.ts[0] < 1800 {
		return 0, false
	}
	t0 := float64(f.ts[0])
	var sx, sy, sxx, sxy float64
	for i := range f.ts {
		x := (float64(f.ts[i]) - t0) / 3600
		y := f.used[i]
		sx += x
		sy += y
		sxx += x * x
		sxy += x * y
	}
	fn := float64(n)
	den := fn*sxx - sx*sx
	if den == 0 {
		return 0, false
	}
	return (fn*sxy - sx*sy) / den, true
}

// FullInHours projects when `avail` bytes are used up at the current rate. ok=false
// when the disk is not filling (or not fast enough to matter: under 1 MB per hour).
func (f *FillTracker) FullInHours(avail uint64) (float64, bool) {
	rate, ok := f.Rate()
	if !ok || rate < 1<<20 {
		return 0, false
	}
	return float64(avail) / rate, true
}
sender.go219 linesraw
package main

import (
	"bytes"
	"compress/gzip"
	"encoding/json"
	"errors"
	"fmt"
	"io"
	"log"
	"net/http"
	"os"
	"path/filepath"
	"sort"
	"strconv"
	"strings"
	"sync"
	"sync/atomic"
	"time"
)

const spoolMaxAge = time.Hour // messages older than this are dropped while offline

// Sender posts reports (gzip JSON) and keeps the ones that could not be delivered in a
// spool (memory plus files in the state directory) for at most an hour. It only ever
// sends: the answer's body is discarded and nothing in it is acted upon.
type Sender struct {
	cfg      Config
	version  string
	client   *http.Client
	mu       sync.Mutex
	flushing sync.Mutex
	// OnResync is called when the server answers 205 Reset Content: its copy of the inventory
	// differs from ours, so the agent sends the full lists once. The status code is the only
	// thing read from an answer.
	OnResync  func()
	sentBytes atomic.Int64
	sentMsgs  atomic.Int64
	queue     []spooled
	lastErr   time.Time
}

type spooled struct {
	at   time.Time
	body []byte // gzip JSON
	file string
}

func NewSender(cfg Config, version string) *Sender {
	return &Sender{cfg: cfg, version: version, client: &http.Client{
		Timeout: 20 * time.Second,
		// Redirects are not followed: the report goes to the configured URL or nowhere.
		CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse },
	}}
}

func (s *Sender) spoolDir() string { return filepath.Join(s.cfg.StateDir, "spool") }

// LoadSpool picks up reports a previous run could not send.
func (s *Sender) LoadSpool() {
	ents, err := os.ReadDir(s.spoolDir())
	if err != nil {
		return
	}
	s.mu.Lock()
	defer s.mu.Unlock()
	for _, e := range ents {
		name := e.Name()
		ts, err := strconv.ParseInt(strings.TrimSuffix(name, ".json.gz"), 10, 64)
		p := filepath.Join(s.spoolDir(), name)
		if err != nil || time.Since(time.Unix(0, ts)) > spoolMaxAge {
			_ = os.Remove(p)
			continue
		}
		b, err := os.ReadFile(p)
		if err != nil {
			continue
		}
		s.queue = append(s.queue, spooled{at: time.Unix(0, ts), body: b, file: p})
	}
	sort.Slice(s.queue, func(i, j int) bool { return s.queue[i].at.Before(s.queue[j].at) })
}

func (s *Sender) Pending() int {
	s.mu.Lock()
	defer s.mu.Unlock()
	return len(s.queue)
}

func Encode(r any) ([]byte, error) {
	raw, err := json.Marshal(r)
	if err != nil {
		return nil, err
	}
	var buf bytes.Buffer
	zw, _ := gzip.NewWriterLevel(&buf, gzip.BestCompression)
	if _, err := zw.Write(raw); err != nil {
		return nil, err
	}
	if err := zw.Close(); err != nil {
		return nil, err
	}
	return buf.Bytes(), nil
}

// Enqueue adds a message and tries to deliver everything queued, oldest first.
func (s *Sender) Enqueue(r any) {
	body, err := Encode(r)
	if err != nil {
		log.Printf("encode report: %v", err)
		return
	}
	now := time.Now()
	item := spooled{at: now, body: body}
	if err := os.MkdirAll(s.spoolDir(), 0o700); err == nil {
		p := filepath.Join(s.spoolDir(), strconv.FormatInt(now.UnixNano(), 10)+".json.gz")
		if os.WriteFile(p, body, 0o600) == nil {
			item.file = p
		}
	}
	s.mu.Lock()
	s.queue = append(s.queue, item)
	s.mu.Unlock()
	go s.Flush()
}

var errRetry = errors.New("retry later")

// retryWaits: one try at once, then after 5 s and 20 s.
var retryWaits = []time.Duration{0, 5 * time.Second, 20 * time.Second}

// Flush sends the queue in order, on its own goroutine (sampling never waits for the
// network). A network error or a 5xx / 429 answer is retried with backoff (5 s, 20 s),
// then left for the next report; anything older than an hour is dropped.
func (s *Sender) Flush() {
	if !s.flushing.TryLock() {
		return // another flush is running and will pick the new report up
	}
	defer s.flushing.Unlock()
	for {
		s.mu.Lock()
		for len(s.queue) > 0 && time.Since(s.queue[0].at) > spoolMaxAge {
			s.drop(s.queue[0])
			s.queue = s.queue[1:]
		}
		if len(s.queue) == 0 {
			s.mu.Unlock()
			return
		}
		item := s.queue[0]
		s.mu.Unlock()
		var err error
		for _, wait := range retryWaits {
			time.Sleep(wait)
			if err = s.post(item.body); err == nil || !errors.Is(err, errRetry) {
				break
			}
		}
		s.mu.Lock()
		if err != nil && errors.Is(err, errRetry) {
			if time.Since(s.lastErr) > 15*time.Minute {
				log.Printf("report not delivered, kept in the spool (%d queued): %v", len(s.queue), err)
				s.lastErr = time.Now()
			}
			s.mu.Unlock()
			return
		}
		if err != nil {
			log.Printf("report dropped: %v", err)
		}
		s.drop(item)
		if len(s.queue) > 0 && s.queue[0].at.Equal(item.at) {
			s.queue = s.queue[1:]
		}
		s.mu.Unlock()
	}
}

func (s *Sender) drop(item spooled) {
	if item.file != "" {
		_ = os.Remove(item.file)
	}
}

func (s *Sender) post(body []byte) error {
	req, err := http.NewRequest(http.MethodPost, s.cfg.Endpoint(), bytes.NewReader(body))
	if err != nil {
		return err
	}
	req.Header.Set("Authorization", "Bearer "+s.cfg.Token)
	req.Header.Set("Content-Type", "application/json")
	req.Header.Set("Content-Encoding", "gzip")
	req.Header.Set("User-Agent", "approvalens-agent/"+s.version)
	res, err := s.client.Do(req)
	if err != nil {
		return fmt.Errorf("%w: %v", errRetry, err)
	}
	_, _ = io.Copy(io.Discard, io.LimitReader(res.Body, 4096))
	res.Body.Close()
	switch {
	case res.StatusCode >= 200 && res.StatusCode < 300:
		s.sentBytes.Add(int64(len(body)))
		s.sentMsgs.Add(1)
		if res.StatusCode == http.StatusResetContent && s.OnResync != nil {
			s.OnResync()
		}
		return nil
	case res.StatusCode == 401 || res.StatusCode == 403:
		return fmt.Errorf("token rejected (HTTP %d): check the token in %s", res.StatusCode, s.cfg.path)
	case res.StatusCode == 413 || res.StatusCode == 400:
		return fmt.Errorf("report refused (HTTP %d)", res.StatusCode)
	default:
		return fmt.Errorf("%w: HTTP %d", errRetry, res.StatusCode)
	}
}

// Sent returns the messages and bytes (compressed, as posted) delivered since start.
func (s *Sender) Sent() (msgs, bytes int64) { return s.sentMsgs.Load(), s.sentBytes.Load() }
signatures.txt129 linesraw
# Heuristic signatures for the Approvalens server agent (embedded into the binary at build time).
# This is NOT an antivirus database: it names a few things that are almost never legitimate on a
# web server, so a match is a reason to look, not proof.
#
#   name <word>   a process whose name (comm, or the base name of argv[0]) equals <word> (case-insensitive)
#   arg  <text>   a process whose command line contains <text> (case-insensitive)
#   busy <word>   a process name allowed to use a whole core for a long time (no "sustained CPU" signal)
#
# Crypto miners and the droppers that start them
name xmrig
name xmrig-notls
name xmr-stak
name xmr-stak-cpu
name xmr-stak-rx
name minerd
name cpuminer
name cpuminer-multi
name cgminer
name bfgminer
name ccminer
name nheqminer
name ethminer
name t-rex
name nbminer
name lolminer
name srbminer
name srbminer-multi
name teamredminer
name gminer
name kdevtmpfsi
name kinsing
name kthreaddi
name kthreaddk
name sysrv
name sysrv-hello
name sysupdate
name networkservice
name sysguard
name dbused
name watchbog
name xmrigdaemon
name xmrigminer
name moneroocean
name c3pool_miner
name tsm
name pnscan
# Command-line fragments of mining software and public pools
arg stratum+tcp://
arg stratum+ssl://
arg stratum+tls://
arg stratum2+tcp://
arg --donate-level
arg --randomx-
arg --cpu-max-threads-hint
arg --coin=monero
arg --algo=rx/0
arg -a rx/0
arg c3pool.com
arg c3pool.org
arg supportxmr.com
arg moneroocean.stream
arg nanopool.org
arg minexmr.com
arg xmrpool.eu
arg hashvault.pro
arg herominers.com
arg 2miners.com
arg f2pool.com
arg unmineable.com
arg monerohash.com
arg xmr.pool.minergate.com
arg pool.hashvault.pro
arg gulf.moneroocean.stream
arg auto.c3pool.org
# Programs that legitimately keep a core busy for a long time
busy mysqld
busy mariadbd
busy postgres
busy mongod
busy redis-server
busy clickhouse-server
busy elasticsearch
busy java
busy ffmpeg
busy x264
busy x265
busy handbrakecli
busy rsync
busy gzip
busy pigz
busy xz
busy zstd
busy bzip2
busy tar
busy borg
busy restic
busy duplicity
busy clamscan
busy clamd
busy freshclam
busy apt
busy apt-get
busy dpkg
busy unattended-upgr
busy yum
busy dnf
busy rpm
busy cc1
busy cc1plus
busy ld
busy go
busy rustc
busy cargo
busy make
busy webpack
busy esbuild
busy imagick
busy convert
busy magick
busy gs
busy ollama
busy updatedb
busy mlocate
busy plocate
busy find
busy fstrim
busy btrfs
busy mdadm
go.mod6 linesraw
module approvalens.com/agent

go 1.22

toolchain go1.23.4
LICENSE22 linesraw
MIT License

Copyright (c) 2026 Kaan Tokalı / Approvalens

Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:

The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.

THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.

Each monitoring subscription includes one server agent. Monitoring