Sunucu ajanı · açık kaynak
Sunucu ajanı, açık açık
Sunucunuzda root olarak bir şey çalıştırmadan önce onu okuyabilmelisiniz. Bu sayfa ajanın tam olarak neleri okuyup gönderdiğini listeler, kaynak kodunu gösterir ve onu kendiniz derleyip bizim dağıttığımız dosyanın aynısını elde ettiğinizi nasıl doğrulayacağınızı anlatır.
- MIT lisansı
- Yalnızca Go standart kütüphanesi, bağımlılık yok
- Statik dosya, yaklaşık 6 MB
- Tekrarlanabilir derleme
Neleri okur, ne zaman gönderir
Ajan 60 saniyede bir yerel olarak ölçüm alır. Bunun çok küçük bir kısmı ağa çıkar: aşağıdaki takvime bakın.
- CPU kullanımı, yük ortalaması, bellek ve swap kullanımı, swap giriş/çıkış hızıNereden okunur:/proc/stat, /proc/loadavg, /proc/meminfo, /proc/vmstatNe zaman gönderilir:5 dakikada bir sinyalde (CPU, bellek); 15 dakikalık özette dakika başına ortalama ve zirve
- Her gerçek dosya sistemi için disk ve inode kullanımı (ext4, xfs, btrfs, zfs…; tmpfs, overlay, snap ve ağ bağlantıları hariç), son 6 saatteki dolma hızı ve “N saatte dolar” tahminiNereden okunur:/proc/self/mountinfo, statfs()Ne zaman gönderilir:En dolu disk ve en yakın tahmin sinyalde; tüm dosya sistemleri özette
- Ağ alma / gönderme hızı (fiziksel arayüzler; lo, docker, veth ve köprüler hariç)Nereden okunur:/proc/net/devNe zaman gönderilir:Özet
- Sunucu adı, işletim sistemi, çekirdek, CPU sayısı, mimari, açılış zamanı, çalışma süresiNereden okunur:/proc/sys/kernel, /etc/os-release, /proc/uptimeNe zaman gönderilir:Özet
- CPU ve belleğe göre ilk 10 süreç: pid, ad, kullanıcı, dosya yolu, 200 karaktere kısaltılmış komut satırı (password=…, --token …, user:pass@ ve mysql -p… maskelenir)Nereden okunur:/proc/PID/stat, cmdline, exeNe zaman gönderilir:Özet
- Dinleyen TCP portları ve UDP servisleri (protokol, port, yalnızca loopback olup olmadığı) ve sahibi olan süreçNereden okunur:/proc/net/tcp, tcp6, udp, udp6, /proc/PID/fdNe zaman gönderilir:Kurulumdan sonra bir kez tam liste; sonra yalnızca eklenen veya çıkanlar
- Aktif crontab satırlarıNereden okunur:/etc/crontab, /etc/cron.d, /var/spool/cronNe zaman gönderilir:Yalnızca değişiklikler
- root'un ve oturum açabilen hesapların SSH anahtarları: tür, SHA-256 parmak izi ve yorum (anahtarın kendisi asla)Nereden okunur:~/.ssh/authorized_keys, authorized_keys2Ne zaman gönderilir:Yalnızca değişiklikler
- setuid / setgid dosyaları (yol, izin, sahip), 6 saatte birNereden okunur:/bin, /sbin, /usr/bin, /usr/sbin, /usr/lib, /usr/libexec, /usr/local, /opt, /tmp, /var/tmp, /dev/shmNe zaman gönderilir:Yalnızca değişiklikler
- Yalnızca web kök dizini ayarlarsanız: yeni, değişmiş ve silinmiş .php, .phtml, .phar, .inc, .js, .mjs, .html, .htm, .shtml, .htaccess ve .user.ini dosyalarının adları, 15 dakikada bir. Değişiklik zamanı ve boyutla karşılaştırılır; içerik asla okunmazNereden okunur:web_roots altında listelediğiniz klasörlerNe zaman gönderilir:Yalnızca değişiklikler (liste başına en fazla 50 ad)
- Günde bir kez: yukarıdaki her listenin özeti (hash), bizdeki kopyanın kaydığını fark edebilmek içinNereden okunur:–Ne zaman gönderilir:Özet, günde bir kez
Ağa neler çıkar
Sinyal (heartbeat) · 5 dakikada bir · yaklaşık 150 bayt
Zaman, ajan sürümü, çalışma süresi, CPU %, bellek %, en dolu disk ve tahmini dolma süresi, şu an tetiklenmiş kurallar. 15 dakika boyunca gelmemesi “sunucu veya ajan sessiz” uyarısıdır.
Özet · 15 dakikada bir · yaklaşık 2 KB
Bir önceki özetten bu yana dakika başına ortalamalar ve zirveler (CPU, bellek, swap, yük, ağ, en dolu disk), dosya sistemleri, ilk süreçler. Grafikler bunu gösterir.
Olay · hemen, dakikada en fazla bir
Yerel bir kural tetiklendiğinde (disk %90 veya 48 saat içinde dolacak, inode %90, 10 dakika boyunca %95 bellek, 10 dakika yoğun swap, 15 dakika boyunca çekirdek sayısının iki katını aşan yük, bir güvenlik sinyali) veya bir liste değiştiğinde (portlar, crontab, SSH anahtarları, setuid, web kök dizinleri), kanıtıyla birlikte.
Envanter · nadiren
Portların, crontab satırlarının, SSH anahtar parmak izlerinin ve setuid dosyalarının tam listesi: kurulumdan sonra bir kez, sonra yalnızca günlük özet kontrolümüz bizdeki kopyanın farklı olduğunu söylerse.
Ölçülen: günde yaklaşık 230 KB veri; HTTPS ek yüküyle (çoğu mesajda bir TLS el sıkışması) kabaca 2 MB. Her şey HTTPS üzerinden gzip JSON'dur: POST https://approvalens.com/api/agent/v1/report, sunucunun anahtarıyla. Bağlantı yokken ajan mesajları en fazla bir saat tutar ve sonra gönderir.
Asla yapmadıkları
Uzaktan komut yok
Hiçbir komut çalıştırmaz, kabuk açmaz ve aldığı hiçbir kodu çalıştırmaz. Bir yanıttan okuduğu tek şey HTTP durum kodudur: 2xx teslim edildi, 205 “tam listelerini yeniden gönder” demektir.
Gelen bağlantı yok
Dinleyen soket açmaz. Tek trafik, ayarlanan adrese giden HTTPS POST isteğidir (düz HTTP yalnızca 127.0.0.1 için kabul edilir).
Dosya içeriği yok
Dosyalarınızın içeriğini asla okumaz. Web kök kontrolleri değişiklik zamanı ve boyutu karşılaştırır ve yalnızca dosya adlarını gönderir. SSH anahtarları parmak izi olarak gönderilir.
Değişiklik yok
Hiçbir şeyi sonlandırmaz, karantinaya almaz veya silmez. Yalnızca kendi durum klasörüne yazar: /var/lib/approvalens-agent (en fazla bir saatlik gönderilmemiş mesajlar, son gönderdikleri, web kök dosya listesi).
Güvenlik sinyalleri sezgiseldir
Bir eşleşme bakmak için bir nedendir, kanıt değildir; eşleşme olmaması da hiçbir şeyi kanıtlamaz: bu bir antivirüs değildir. Sinyaller şunlardır: kripto madenci listesinden (aşağıdaki signatures.txt) bir süreç adı veya argümanı; /tmp, /var/tmp veya /dev/shm altındaki bir çalıştırılabilir dosya; başladıktan sonra silinmiş (yükseltmelerin meşru olarak yaptığı sistem klasörleri dışında) veya bellekten çalışan (memfd) bir dosya; 15 dakika boyunca tam bir çekirdek kullanan bir süreç (veritabanları, sıkıştırıcılar, derleyiciler ve yedekleme araçları hariç); ve ilk 24 saatten (temel durum) sonra yukarıdaki listelerde beliren her yeni şey. Bir sinyali onaylayabilirsiniz; temel duruma eklenir.
Yetkiler
Varsayılan olarak hiçbir yetkisi olmadan kendi kullanıcısıyla, approvalens-agent olarak çalışır. Linux bu durumda birkaç şeyi ondan gizler ve hesabınız hangilerini gösterir: diğer kullanıcıların süreçlerinin dosya yolu (komut satırına geri döner), diğer kullanıcıların dinleyen soketlerinin sahibi, kullanıcı crontab'ları ve diğer kullanıcıların SSH anahtarları. Bunları da eklemek için --root ile kurun: servis root olarak çalışır ama yalnızca CAP_DAC_READ_SEARCH ve CAP_SYS_PTRACE yetkilerini (bu dosyaları ve /proc/PID/exe ile fd'yi okumak için) salt okunur bir dosya sisteminde tutar.
Kaynak sınırları ve systemd sıkılaştırması
Hedef: ortalamada bir çekirdeğin %1'inden az ve 30 MB'tan az bellek. Servis sert sınırlar uygular ve süreci kilitler:
[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
Bir test makinesinde ölçülen: ortalama bir çekirdeğin %0.12 kadarı, 12 MB bellek.
Kurulum
Anahtarını almak için hesabınızda bir sunucu oluşturun, sonra sunucuda: /account
curl -fsSL https://approvalens.com/agent/install.sh | sudo bash -s -- --token <TOKEN>
Önce hiçbir şeyi değiştirmeden her adımı görün:
curl -fsSL https://approvalens.com/agent/install.sh | sudo bash -s -- --token <TOKEN> --dry-run
Kendi derlediğiniz dosyayı kurun:
curl -fsSLO https://approvalens.com/agent/install.sh sudo bash install.sh --token <TOKEN> --binary ./approvalens-agent-linux-amd64
Göndermeden, tam olarak ne göndereceğini yazdırın:
sudo -u approvalens-agent approvalens-agent check --config /etc/approvalens-agent.yaml
Kaldırma
Servisi durdurur; dosyayı, ayarı, durum klasörünü, servis tanımını ve kullanıcıyı kaldırır:
curl -fsSL https://approvalens.com/agent/install.sh | sudo bash -s -- --uninstall
Kendiniz derleyin ve karşılaştırın
Derlemeler tekrarlanabilir: aynı Go sürümü (go.mod içindeki toolchain satırı; Go bunu kendisi indirir) her makinede bayt bayt aynı dosyaları üretir. Kaynak arşivinden derleyin ve SHA-256 özetini bizimkiyle karşılaştırın:
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Özetler SHA256SUMS içindeki satırlarla aynı olmalı. Değilse kurmayın ve bize haber verin.
İndirmeler · sürüm 0.1.0
- Kaynak kod (tar.gz) approvalens-agent-0.1.0-src.tar.gz
- SHA256SUMS SHA256SUMS
- Linux x86-64 dosyası approvalens-agent-linux-amd64
- Linux ARM64 dosyası approvalens-agent-linux-arm64
- install.sh
8e6ab0397de8706f17304fe0392a284e5192a2b9886efc3315f48dff645b8539 approvalens-agent-linux-amd64 5d7ba23c0572174cc0135bd1fe9c7e0d9b04809ea7c462129bd88f089f3a0d49 approvalens-agent-linux-arm64 0b0a10618fa7916d47357cacee3091bb704f44a717f9b4139292b6ff176aa1db approvalens-agent-0.1.0-src.tar.gz
Kaynak kod
0.1.0 sürümünün ana dosyaları, derlendiği haliyle. Yukarıdaki arşivde testler ve test verileri de var.
main.go225 satırraw
// 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 satırraw
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 satırraw
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 satırraw
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 satırraw
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 satırraw
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 satırraw
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 satırraw
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 satırraw
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 satırraw
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 satırraw
# 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 satırraw
module approvalens.com/agent
go 1.22
toolchain go1.23.4
LICENSE22 satırraw
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.
Her izleme aboneliği bir sunucu ajanı içerir. İzleme