statuspage: monitor lab/home machines with a small status page

Go backend probes machines over ICMP/TCP, runs a configurable script over
SSH for checks and numeric metrics, and keeps history in SQLite. Preact +
Vite frontend shows machine groups, per-machine metric plots, and
cross-machine metric aggregates.

- server/: status, history, metrics, aggregate, refresh APIs
- src/: Preact + Vite frontend
- Single YAML config; documented example in example.config.yaml
- SQLite data stored under data.volume/ (bind-mounted in Docker)
- Multi-stage alpine Dockerfile with BuildKit cache mounts
- LICENSE: AGPL-3.0
This commit is contained in:
2026-08-01 02:27:47 +02:00
commit 95eab7f584
21 changed files with 5208 additions and 0 deletions
+197
View File
@@ -0,0 +1,197 @@
package main
import (
"fmt"
"os"
"time"
"gopkg.in/yaml.v3"
)
type Duration struct {
time.Duration
}
func (d *Duration) UnmarshalYAML(node *yaml.Node) error {
var s string
if err := node.Decode(&s); err != nil {
return err
}
dur, err := time.ParseDuration(s)
if err != nil {
return fmt.Errorf("invalid duration %q: %w", s, err)
}
d.Duration = dur
return nil
}
type PingConfig struct {
Interval Duration `yaml:"interval"`
Timeout Duration `yaml:"timeout"`
TCPPort int `yaml:"tcp_port"`
}
type SSHConfig struct {
Interval Duration `yaml:"interval"`
User string `yaml:"user"`
Port int `yaml:"port"`
Key string `yaml:"key"`
Script string `yaml:"script"`
}
type MachineConfig struct {
Name string `yaml:"name"`
Host string `yaml:"host"`
Ping *PingConfig `yaml:"ping"`
SSH *SSHConfig `yaml:"ssh"`
}
type GroupConfig struct {
Name string `yaml:"name"`
Machines []MachineConfig `yaml:"machines"`
}
type MetricRange struct {
Min *float64 `yaml:"min"`
Max *float64 `yaml:"max"`
}
type RawConfig struct {
Title string `yaml:"title"`
Interactive bool `yaml:"interactive"`
SharedMetricWindow bool `yaml:"shared_metric_window"`
Ping *PingConfig `yaml:"ping"`
SSH *SSHConfig `yaml:"ssh"`
Metrics map[string]MetricRange `yaml:"metrics"`
Groups []GroupConfig `yaml:"groups"`
Machines []MachineConfig `yaml:"machines"`
}
type Machine struct {
ID string
Name string
Host string
Group string
Ping PingConfig
SSH SSHConfig
}
type Config struct {
Title string
Interactive bool
SharedMetricWindow bool
Ping PingConfig
SSH SSHConfig
Metrics map[string]MetricRange
Groups []GroupConfig
Machines []Machine
}
func LoadConfig(path string) (*Config, error) {
data, err := os.ReadFile(path)
if err != nil {
return nil, err
}
var raw RawConfig
if err := yaml.Unmarshal(data, &raw); err != nil {
return nil, fmt.Errorf("parsing %s: %w", path, err)
}
cfg := &Config{
Title: raw.Title,
Interactive: raw.Interactive,
SharedMetricWindow: raw.SharedMetricWindow,
Ping: PingConfig{
Interval: Duration{time.Second * 10},
Timeout: Duration{time.Second * 2},
TCPPort: 22,
},
SSH: SSHConfig{
Interval: Duration{time.Minute * 10},
User: "root",
Port: 22,
},
}
if raw.Ping != nil {
if raw.Ping.Interval.Duration > 0 {
cfg.Ping.Interval = raw.Ping.Interval
}
if raw.Ping.Timeout.Duration > 0 {
cfg.Ping.Timeout = raw.Ping.Timeout
}
if raw.Ping.TCPPort != 0 {
cfg.Ping.TCPPort = raw.Ping.TCPPort
}
}
if raw.SSH != nil {
if raw.SSH.Interval.Duration > 0 {
cfg.SSH.Interval = raw.SSH.Interval
}
if raw.SSH.User != "" {
cfg.SSH.User = raw.SSH.User
}
if raw.SSH.Port != 0 {
cfg.SSH.Port = raw.SSH.Port
}
cfg.SSH.Key = raw.SSH.Key
cfg.SSH.Script = raw.SSH.Script
}
if cfg.Title == "" {
cfg.Title = "Status Page"
}
cfg.Groups = raw.Groups
cfg.Metrics = raw.Metrics
for _, g := range raw.Groups {
for _, mc := range g.Machines {
cfg.Machines = append(cfg.Machines, resolveMachine(mc, g.Name, cfg))
}
}
for _, mc := range raw.Machines {
cfg.Machines = append(cfg.Machines, resolveMachine(mc, "", cfg))
}
return cfg, nil
}
func resolveMachine(mc MachineConfig, group string, cfg *Config) Machine {
m := Machine{
ID: mc.Host,
Name: mc.Host,
Host: mc.Host,
Group: group,
Ping: cfg.Ping,
SSH: cfg.SSH,
}
if mc.Name != "" {
m.Name = mc.Name
}
if mc.Ping != nil {
if mc.Ping.Interval.Duration > 0 {
m.Ping.Interval = mc.Ping.Interval
}
if mc.Ping.Timeout.Duration > 0 {
m.Ping.Timeout = mc.Ping.Timeout
}
if mc.Ping.TCPPort != 0 {
m.Ping.TCPPort = mc.Ping.TCPPort
}
}
if mc.SSH != nil {
if mc.SSH.Interval.Duration > 0 {
m.SSH.Interval = mc.SSH.Interval
}
if mc.SSH.User != "" {
m.SSH.User = mc.SSH.User
}
if mc.SSH.Port != 0 {
m.SSH.Port = mc.SSH.Port
}
if mc.SSH.Key != "" {
m.SSH.Key = mc.SSH.Key
}
if mc.SSH.Script != "" {
m.SSH.Script = mc.SSH.Script
}
}
return m
}
+402
View File
@@ -0,0 +1,402 @@
package main
import (
"context"
"database/sql"
"log"
"os"
"sort"
"time"
_ "github.com/mattn/go-sqlite3"
)
const historyPerMachine = 5000
const historyTotalCap = 200000
const metricsPerName = 5000
const historyPruneInterval = 10 * time.Minute
type HistoryEntry struct {
TS int64 `json:"ts"`
Status string `json:"status"`
IP string `json:"ip"`
}
type MetricEntry struct {
TS int64 `json:"ts"`
Name string `json:"name"`
Value float64 `json:"value"`
Unit string `json:"unit"`
}
type History struct {
db *sql.DB
path string
}
func NewHistory(path string) (*History, error) {
dsn := "file:" + path +
"?_journal_mode=WAL" +
"&_synchronous=NORMAL" +
"&_busy_timeout=5000" +
"&_foreign_keys=on" +
"&_txlock=immediate"
db, err := sql.Open("sqlite3", dsn)
if err != nil {
return nil, err
}
db.SetMaxOpenConns(1)
const schema = `
CREATE TABLE IF NOT EXISTS history (
id INTEGER PRIMARY KEY AUTOINCREMENT,
machine TEXT NOT NULL,
ts INTEGER NOT NULL,
status TEXT NOT NULL,
ip TEXT NOT NULL DEFAULT ''
);
CREATE INDEX IF NOT EXISTS idx_history_machine_ts ON history(machine, ts);
CREATE TABLE IF NOT EXISTS metrics (
id INTEGER PRIMARY KEY AUTOINCREMENT,
machine TEXT NOT NULL,
name TEXT NOT NULL,
ts INTEGER NOT NULL,
value REAL NOT NULL,
unit TEXT NOT NULL DEFAULT ''
);
CREATE INDEX IF NOT EXISTS idx_metrics_machine_name_ts ON metrics(machine, name, ts);`
if _, err := db.Exec(schema); err != nil {
db.Close()
return nil, err
}
return &History{db: db, path: path}, nil
}
func (h *History) Close() error {
return h.db.Close()
}
func (h *History) Size() int64 {
if fi, err := os.Stat(h.path); err == nil {
return fi.Size()
}
return 0
}
func (h *History) Stats() (rows int, firstTS int64) {
h.db.QueryRow(`SELECT COUNT(*), COALESCE(MIN(ts), 0) FROM history`).Scan(&rows, &firstTS)
return
}
func (h *History) Record(machine, status, ip string) {
_, err := h.db.Exec(
`INSERT INTO history (machine, ts, status, ip) VALUES (?, ?, ?, ?)`,
machine, time.Now().Unix(), status, ip,
)
if err != nil {
log.Printf("history insert (%s): %v", machine, err)
}
}
func (h *History) RecordMetric(machine, name string, ts int64, value float64, unit string) {
_, err := h.db.Exec(
`INSERT INTO metrics (machine, name, ts, value, unit) VALUES (?, ?, ?, ?, ?)`,
machine, name, ts, value, unit,
)
if err != nil {
log.Printf("metric insert (%s/%s): %v", machine, name, err)
}
}
func (h *History) Query(machine string, limit int) ([]HistoryEntry, error) {
if limit <= 0 || limit > historyPerMachine {
limit = historyPerMachine
}
rows, err := h.db.Query(
`SELECT ts, status, ip FROM history WHERE machine = ? ORDER BY id DESC LIMIT ?`,
machine, limit,
)
if err != nil {
return nil, err
}
defer rows.Close()
entries := []HistoryEntry{}
for rows.Next() {
var e HistoryEntry
if err := rows.Scan(&e.TS, &e.Status, &e.IP); err != nil {
return nil, err
}
entries = append(entries, e)
}
return entries, rows.Err()
}
func (h *History) QueryMetrics(machine, name string, minTS, maxTS int64, maxPoints int, shared bool) ([]MetricEntry, error) {
if maxPoints <= 0 || maxPoints > metricsPerName*5 {
maxPoints = 100
}
q := `SELECT ts, name, value, unit FROM metrics WHERE machine = ?`
args := []any{machine}
if name != "" {
q += ` AND name = ?`
args = append(args, name)
}
if minTS > 0 {
q += ` AND ts >= ?`
args = append(args, minTS)
}
if maxTS > 0 {
q += ` AND ts < ?`
args = append(args, maxTS)
}
q += ` ORDER BY id ASC LIMIT ?`
args = append(args, metricsPerName*5)
rows, err := h.db.Query(q, args...)
if err != nil {
return nil, err
}
defer rows.Close()
entries := []MetricEntry{}
for rows.Next() {
var e MetricEntry
if err := rows.Scan(&e.TS, &e.Name, &e.Value, &e.Unit); err != nil {
return nil, err
}
entries = append(entries, e)
}
if err := rows.Err(); err != nil {
return nil, err
}
var lo, hi int64
if shared {
lo = minTS
hi = maxTS
}
return downsampleMetrics(entries, lo, hi, maxPoints), nil
}
func downsampleMetrics(entries []MetricEntry, lo, hi int64, maxPoints int) []MetricEntry {
if maxPoints < 1 {
return entries
}
groups := map[string][]MetricEntry{}
order := []string{}
for _, e := range entries {
if _, ok := groups[e.Name]; !ok {
order = append(order, e.Name)
}
groups[e.Name] = append(groups[e.Name], e)
}
out := []MetricEntry{}
for _, name := range order {
out = append(out, downsampleGroup(groups[name], lo, hi, maxPoints)...)
}
sort.Slice(out, func(i, j int) bool { return out[i].TS < out[j].TS })
return out
}
func downsampleGroup(g []MetricEntry, lo, hi int64, maxPoints int) []MetricEntry {
if len(g) <= maxPoints {
return g
}
return alignGroup(g, lo, hi, maxPoints)
}
// alignGroup downsamples g to at most maxPoints buckets whose boundaries are
// aligned to the [lo, hi] window, so all machines share the same bucket centers.
func alignGroup(g []MetricEntry, lo, hi int64, maxPoints int) []MetricEntry {
if len(g) == 0 {
return nil
}
if lo <= 0 {
lo = g[0].TS
}
if hi <= 0 {
hi = g[len(g)-1].TS
}
span := hi - lo
if span <= 0 {
span = 1
}
bucketSize := float64(span) / float64(maxPoints)
buckets := make([][]MetricEntry, maxPoints)
for _, e := range g {
idx := int(float64(e.TS-lo) / bucketSize)
if idx >= maxPoints {
idx = maxPoints - 1
}
if idx < 0 {
idx = 0
}
buckets[idx] = append(buckets[idx], e)
}
res := make([]MetricEntry, 0, maxPoints)
for i, b := range buckets {
if len(b) == 0 {
continue
}
var sum float64
for _, e := range b {
sum += e.Value
}
last := b[len(b)-1]
res = append(res, MetricEntry{
TS: lo + int64((float64(i)+0.5)*bucketSize),
Name: last.Name,
Value: sum / float64(len(b)),
Unit: last.Unit,
})
}
return res
}
func (h *History) QueryMetricsAll(minTS, maxTS int64) ([]MetricAggregate, [2]int64, error) {
q := `SELECT machine, ts, name, value, unit FROM metrics`
args := []any{}
if minTS > 0 {
q += ` WHERE ts >= ?`
args = append(args, minTS)
}
if maxTS > 0 {
if len(args) > 0 {
q += ` AND ts < ?`
} else {
q += ` WHERE ts < ?`
}
args = append(args, maxTS)
}
q += ` ORDER BY machine, name, id ASC LIMIT ?`
args = append(args, metricsPerName*10)
rows, err := h.db.Query(q, args...)
if err != nil {
return nil, [2]int64{}, err
}
defer rows.Close()
type rawEntry struct {
machine string
entry MetricEntry
}
raw := []rawEntry{}
var firstTS, lastTS int64
for rows.Next() {
var e rawEntry
if err := rows.Scan(&e.machine, &e.entry.TS, &e.entry.Name, &e.entry.Value, &e.entry.Unit); err != nil {
return nil, [2]int64{}, err
}
if firstTS == 0 || e.entry.TS < firstTS {
firstTS = e.entry.TS
}
if e.entry.TS > lastTS {
lastTS = e.entry.TS
}
raw = append(raw, e)
}
if err := rows.Err(); err != nil {
return nil, [2]int64{}, err
}
if minTS > 0 && firstTS < minTS {
firstTS = minTS
}
if maxTS > 0 && lastTS > maxTS {
lastTS = maxTS
}
type nameGroup struct {
unit string
series []MetricSeries
}
byName := map[string]*nameGroup{}
order := []string{}
for _, r := range raw {
ng, ok := byName[r.entry.Name]
if !ok {
ng = &nameGroup{unit: r.entry.Unit}
byName[r.entry.Name] = ng
order = append(order, r.entry.Name)
}
found := false
for i := range ng.series {
if ng.series[i].Machine == r.machine {
ng.series[i].Samples = append(ng.series[i].Samples, r.entry)
found = true
break
}
}
if !found {
ng.series = append(ng.series, MetricSeries{Machine: r.machine, Samples: []MetricEntry{r.entry}})
}
}
out := make([]MetricAggregate, 0, len(order))
for _, name := range order {
ng := byName[name]
agg := MetricAggregate{Name: name, Unit: ng.unit, Series: make([]MetricSeries, 0, len(ng.series))}
for _, s := range ng.series {
sort.Slice(s.Samples, func(i, j int) bool { return s.Samples[i].TS < s.Samples[j].TS })
agg.Series = append(agg.Series, s)
}
out = append(out, agg)
}
return out, [2]int64{firstTS, lastTS}, nil
}
type MetricSeries struct {
Machine string `json:"machine"`
Samples []MetricEntry `json:"samples"`
}
type MetricAggregate struct {
Name string `json:"name"`
Unit string `json:"unit"`
Series []MetricSeries `json:"series"`
}
func (h *History) prune() {
res, err := h.db.Exec(
`DELETE FROM history WHERE id NOT IN (SELECT id FROM history ORDER BY id DESC LIMIT ?)`,
historyTotalCap,
)
if err != nil {
log.Printf("history prune: %v", err)
} else {
n, _ := res.RowsAffected()
if n > 0 {
log.Printf("history pruned %d rows (kept last %d)", n, historyTotalCap)
}
}
res, err = h.db.Exec(`
DELETE FROM metrics WHERE id NOT IN (
SELECT id FROM (
SELECT id,
ROW_NUMBER() OVER (PARTITION BY machine, name ORDER BY id DESC) AS rn
FROM metrics
) WHERE rn <= ?
)`, metricsPerName)
if err != nil {
log.Printf("metrics prune: %v", err)
return
}
n, _ := res.RowsAffected()
if n > 0 {
log.Printf("metrics pruned %d rows (kept last %d per name)", n, metricsPerName)
}
}
func (h *History) pruneLoop(ctx context.Context) {
ticker := time.NewTicker(historyPruneInterval)
defer ticker.Stop()
h.prune()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
h.prune()
}
}
}
+98
View File
@@ -0,0 +1,98 @@
package main
import (
"context"
"log"
"net/http"
"os"
"os/exec"
"path/filepath"
"strings"
"github.com/spf13/pflag"
)
func buildFrontend() error {
cmd := "npm"
if _, err := exec.LookPath("npm"); err != nil {
cmd = "bun"
}
c := exec.Command(cmd, "run", "build")
c.Stdout = os.Stdout
c.Stderr = os.Stderr
log.Printf("running %s run build", cmd)
return c.Run()
}
func spaHandler() http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
if strings.HasPrefix(r.URL.Path, "/api/") {
http.NotFound(w, r)
return
}
rel := filepath.Clean(strings.TrimPrefix(r.URL.Path, "/"))
if rel != "." && rel != "" {
p := filepath.Join("dist", rel)
if fi, err := os.Stat(p); err == nil && !fi.IsDir() {
http.ServeFile(w, r, p)
return
}
}
http.ServeFile(w, r, filepath.Join("dist", "index.html"))
}
}
func main() {
dev := pflag.Bool("dev", false, "run `npm run build` before serving dist/")
devServer := pflag.Bool("dev-server", false, "frontend served by the vite dev server, do not serve dist/")
addr := pflag.String("addr", ":5000", "listen address")
configPath := pflag.StringP("config", "c", "config.local.yaml", "yaml config file")
dbPath := pflag.String("db", "data.volume/history.db", "sqlite history database file")
pflag.Parse()
if dir := filepath.Dir(*dbPath); dir != "." {
if err := os.MkdirAll(dir, 0o755); err != nil {
log.Fatalf("creating db dir: %v", err)
}
}
if *dev {
if err := buildFrontend(); err != nil {
log.Fatalf("build failed: %v", err)
}
}
cfg, err := LoadConfig(*configPath)
if err != nil {
log.Fatalf("loading config: %v", err)
}
log.Printf("loaded %d machines from %s", len(cfg.Machines), *configPath)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
history, err := NewHistory(*dbPath)
if err != nil {
log.Fatalf("opening history db: %v", err)
}
defer history.Close()
go history.pruneLoop(ctx)
mon := NewMonitor(cfg, history)
mon.Start(ctx)
mux := http.NewServeMux()
mux.HandleFunc("GET /api/status", mon.statusHandler)
mux.HandleFunc("GET /api/history/{machine}", mon.historyHandler)
mux.HandleFunc("GET /api/metrics/aggregate", mon.metricsAggregateHandler)
mux.HandleFunc("GET /api/metrics/{machine}", mon.metricsHandler)
mux.HandleFunc("POST /api/refresh/{machine}", mon.refreshHandler)
if !*devServer {
mux.HandleFunc("GET /", spaHandler())
}
log.Printf("listening on %s", *addr)
if err := http.ListenAndServe(*addr, mux); err != nil {
log.Fatal(err)
}
}
+771
View File
@@ -0,0 +1,771 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"math/rand/v2"
"net"
"net/http"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
"golang.org/x/crypto/ssh"
"golang.org/x/net/icmp"
"golang.org/x/net/ipv4"
)
const statusUp = "up"
const statusDegraded = "degraded"
const statusDown = "down"
const statusUnknown = "unknown"
const checkOK = "ok"
const checkFail = "fail"
const checkNA = "na"
const maxIPHistory = 10
type sshResultState struct {
Ok bool `json:"ok"`
Result string `json:"result"`
ExitCode *int `json:"exitCode"`
LastRun time.Time `json:"lastRun"`
Error string `json:"error"`
}
type metricState struct {
value float64
unit string
updated time.Time
}
type machineState struct {
mu sync.RWMutex
status string
icmp string
tcp string
lastPing time.Time
ip string
ips []string
sshResult *sshResultState
metrics map[string]metricState
}
type metricStatus struct {
Name string `json:"name"`
Value float64 `json:"value"`
Unit string `json:"unit"`
Updated *time.Time `json:"updated"`
}
type machineStatus struct {
ID string `json:"id"`
Name string `json:"name"`
Host string `json:"host"`
IP string `json:"ip"`
IPs []string `json:"ips"`
Status string `json:"status"`
ICMP string `json:"icmp"`
TCP string `json:"tcp"`
LastPing *time.Time `json:"lastPing"`
SSHConfigured bool `json:"sshConfigured"`
SSH *sshResultState `json:"ssh"`
Metrics []metricStatus `json:"metrics"`
}
type groupStatus struct {
Name string `json:"name"`
Machines []machineStatus `json:"machines"`
}
type statsPayload struct {
PingInterval string `json:"pingInterval"`
PingTimeout string `json:"pingTimeout"`
TCPPort int `json:"tcpPort"`
SSHInterval string `json:"sshInterval"`
SSHEnabled bool `json:"sshEnabled"`
DBSize int64 `json:"dbSize"`
DBRows int `json:"dbRows"`
DBSizeMonth int64 `json:"dbSizeMonth"`
}
type metricRange struct {
Min *float64 `json:"min,omitempty"`
Max *float64 `json:"max,omitempty"`
}
type statusPayload struct {
Title string `json:"title"`
Interactive bool `json:"interactive"`
SharedMetricWindow bool `json:"sharedMetricWindow"`
Stats statsPayload `json:"stats"`
MetricRanges map[string]metricRange `json:"metricRanges"`
Groups []groupStatus `json:"groups"`
Machines []machineStatus `json:"machines"`
}
type Monitor struct {
cfg *Config
machines []Machine
machinesByID map[string]Machine
interactive bool
ctx context.Context
history *History
icmpMu sync.Mutex
icmp *icmp.PacketConn
icmpSeq atomic.Uint32
icmpID int
statesMu sync.RWMutex
states map[string]*machineState
}
func NewMonitor(cfg *Config, history *History) *Monitor {
m := &Monitor{
cfg: cfg,
machines: cfg.Machines,
machinesByID: make(map[string]Machine, len(cfg.Machines)),
interactive: cfg.Interactive,
ctx: context.Background(),
history: history,
states: make(map[string]*machineState, len(cfg.Machines)),
}
for _, mc := range cfg.Machines {
m.states[mc.ID] = &machineState{status: statusUnknown, icmp: checkNA, tcp: checkNA}
m.machinesByID[mc.ID] = mc
}
conn, err := icmp.ListenPacket("udp4", "0.0.0.0")
if err != nil {
conn, err = icmp.ListenPacket("ip4:icmp", "0.0.0.0")
}
if err != nil {
log.Printf("icmp unavailable (%v); falling back to TCP-only probes", err)
m.icmp = nil
} else {
m.icmp = conn
if la, ok := conn.LocalAddr().(*net.UDPAddr); ok {
m.icmpID = la.Port
}
log.Printf("icmp ok, will use it with tcp fallback")
}
return m
}
func (m *Monitor) Start(ctx context.Context) {
m.ctx = ctx
for _, mc := range m.machines {
go m.pingLoop(ctx, mc)
go m.sshLoop(ctx, mc)
}
}
func (m *Monitor) Refresh(mc Machine) {
go m.probe(m.ctx, mc)
go m.runSSH(m.ctx, mc)
}
func (m *Monitor) pingLoop(ctx context.Context, mc Machine) {
run := func() {
m.probe(ctx, mc)
}
run()
ticker := time.NewTicker(mc.Ping.Interval.Duration)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
run()
}
}
}
func (m *Monitor) state(mc Machine) *machineState {
m.statesMu.RLock()
s := m.states[mc.ID]
m.statesMu.RUnlock()
return s
}
func (m *Monitor) probe(ctx context.Context, mc Machine) {
st := m.state(mc)
if st == nil {
return
}
timeout := mc.Ping.Timeout.Duration
if timeout <= 0 {
timeout = 2 * time.Second
}
resolveCtx, cancel := context.WithTimeout(ctx, timeout)
ips, err := net.DefaultResolver.LookupHost(resolveCtx, mc.Host)
cancel()
if err != nil {
m.setProbeResult(st, statusDown, checkFail, checkFail, "", nil)
m.recordHistory(mc.ID, statusDown, "")
log.Printf("ping %s -> %s (resolve failed: %v)", mc.Host, statusDown, err)
return
}
// only ping over ipv4, ignore ipv6 addresses
var v4 []string
for _, a := range ips {
if net.ParseIP(a).To4() != nil {
v4 = append(v4, a)
}
}
if len(v4) == 0 {
m.setProbeResult(st, statusDown, checkFail, checkFail, "", nil)
m.recordHistory(mc.ID, statusDown, "")
log.Printf("ping %s -> %s (no ipv4 addresses)", mc.Host, statusDown)
return
}
ip := v4[0]
m.recordIP(st, v4)
icmpOK := m.icmpPing(net.ParseIP(ip), timeout)
tcpOK := m.tcpPing(net.ParseIP(ip), mc.Ping.TCPPort, timeout)
icmp := checkNA
if m.icmp != nil {
if icmpOK {
icmp = checkOK
} else {
icmp = checkFail
}
}
tcp := checkNA
if tcpOK {
tcp = checkOK
} else {
tcp = checkFail
}
status := statusDown
if icmpOK || tcpOK {
status = statusUp
st.mu.RLock()
sshBad := st.sshResult != nil && !st.sshResult.Ok
st.mu.RUnlock()
if (icmp == checkOK && tcp == checkFail) ||
(icmp == checkFail && tcp == checkOK) || sshBad {
status = statusDegraded
}
}
now := time.Now()
m.setProbeResult(st, status, icmp, tcp, ip, &now)
m.recordHistory(mc.ID, status, ip)
log.Printf("ping %s (%s) icmp=%s tcp=%s -> %s", mc.Host, ip, icmp, tcp, status)
}
func (m *Monitor) recordHistory(machine, status, ip string) {
if m.history != nil {
m.history.Record(machine, status, ip)
}
}
func (m *Monitor) setProbeResult(st *machineState, status, icmp, tcp, ip string, last *time.Time) {
st.mu.Lock()
defer st.mu.Unlock()
st.status = status
st.icmp = icmp
st.tcp = tcp
if ip != "" {
st.ip = ip
}
if last != nil {
st.lastPing = *last
}
}
func (m *Monitor) recordIP(st *machineState, ips []string) {
st.mu.Lock()
defer st.mu.Unlock()
for _, ip := range ips {
if ip == st.ip {
continue
}
st.ips = append(st.ips, ip)
}
if len(st.ips) > maxIPHistory {
st.ips = st.ips[len(st.ips)-maxIPHistory:]
}
}
func (m *Monitor) icmpPing(ip net.IP, timeout time.Duration) bool {
if m.icmp == nil || ip == nil {
return false
}
m.icmpMu.Lock()
defer m.icmpMu.Unlock()
seq := int(m.icmpSeq.Add(1))
wm := icmp.Message{
Type: ipv4.ICMPTypeEcho,
Code: 0,
Body: &icmp.Echo{ID: m.icmpID, Seq: seq, Data: []byte("lab-status")},
}
wb, err := wm.Marshal(nil)
if err != nil {
return false
}
if _, err := m.icmp.WriteTo(wb, &net.UDPAddr{IP: ip}); err != nil {
return false
}
m.icmp.SetReadDeadline(time.Now().Add(timeout))
reply := make([]byte, 1500)
for {
n, peer, err := m.icmp.ReadFrom(reply)
if err != nil {
return false
}
pa, ok := peer.(*net.UDPAddr)
if !ok || !pa.IP.Equal(ip) {
continue
}
rm, err := icmp.ParseMessage(1, reply[:n])
if err != nil {
continue
}
if rm.Type != ipv4.ICMPTypeEchoReply {
continue
}
if echo, ok := rm.Body.(*icmp.Echo); ok && echo.ID == m.icmpID && echo.Seq == seq {
return true
}
}
}
func (m *Monitor) tcpPing(ip net.IP, port int, timeout time.Duration) bool {
if ip == nil {
return false
}
addr := net.JoinHostPort(ip.String(), strconv.Itoa(port))
conn, err := net.DialTimeout("tcp", addr, timeout)
if err != nil {
return false
}
conn.Close()
return true
}
func (m *Monitor) sshLoop(ctx context.Context, mc Machine) {
if mc.SSH.Key == "" || mc.SSH.Script == "" {
return
}
interval := mc.SSH.Interval.Duration
if interval <= 0 {
interval = 10 * time.Minute
}
jitter := interval / 4
if jitter <= 0 {
jitter = interval
}
// stagger the first run within [0, interval) so machines don't all connect at once
timer := time.NewTimer(rand.N(interval))
defer timer.Stop()
for {
select {
case <-ctx.Done():
return
case <-timer.C:
m.runSSH(ctx, mc)
// drift each run by a random offset to keep machines out of phase
timer.Reset(interval + rand.N(jitter))
}
}
}
func (m *Monitor) runSSH(ctx context.Context, mc Machine) {
start := time.Now()
st := m.state(mc)
if st == nil {
return
}
key, err := ssh.ParsePrivateKey([]byte(mc.SSH.Key))
if err != nil {
log.Printf("ssh %s: bad key: %v", mc.Host, err)
m.setSSHError(st, fmt.Sprintf("bad ssh key: %v", err))
return
}
cfg := &ssh.ClientConfig{
User: mc.SSH.User,
Auth: []ssh.AuthMethod{ssh.PublicKeys(key)},
HostKeyCallback: ssh.InsecureIgnoreHostKey(),
Timeout: mc.Ping.Timeout.Duration,
}
addr := net.JoinHostPort(mc.Host, strconv.Itoa(mc.SSH.Port))
log.Printf("ssh %s: connecting to %s", mc.Host, addr)
client, err := ssh.Dial("tcp", addr, cfg)
if err != nil {
log.Printf("ssh %s: dial failed: %v", mc.Host, err)
m.setSSHError(st, fmt.Sprintf("ssh dial: %v", err))
return
}
defer client.Close()
log.Printf("ssh %s: connected", mc.Host)
session, err := client.NewSession()
if err != nil {
log.Printf("ssh %s: session failed: %v", mc.Host, err)
m.setSSHError(st, fmt.Sprintf("ssh session: %v", err))
return
}
defer session.Close()
timeout := mc.SSH.Interval.Duration
if timeout > 30*time.Minute {
timeout = 30 * time.Minute
}
log.Printf("ssh %s: running script", mc.Host)
type cmdResult struct {
out []byte
err error
}
done := make(chan cmdResult, 1)
go func() {
out, err := session.CombinedOutput(mc.SSH.Script)
done <- cmdResult{out, err}
}()
select {
case r := <-done:
var exitCode *int
if r.err != nil {
var ee *ssh.ExitError
if errors.As(r.err, &ee) {
ec := ee.ExitStatus()
exitCode = &ec
log.Printf("ssh %s: script done in %s, exit %d", mc.Host, time.Since(start), ec)
} else {
log.Printf("ssh %s: script failed: %v", mc.Host, r.err)
m.setSSHError(st, fmt.Sprintf("ssh command: %v", r.err))
return
}
} else {
log.Printf("ssh %s: script done in %s, exit 0", mc.Host, time.Since(start))
}
st.mu.Lock()
st.sshResult = &sshResultState{
Ok: exitCode == nil,
Result: string(r.out),
ExitCode: exitCode,
LastRun: time.Now(),
}
st.mu.Unlock()
if exitCode == nil {
m.accumulateMetrics(st, mc.ID, string(r.out), time.Now())
}
case <-time.After(timeout):
session.Close()
log.Printf("ssh %s: timed out after %s", mc.Host, timeout)
m.setSSHError(st, fmt.Sprintf("ssh command timed out after %s", timeout))
case <-ctx.Done():
log.Printf("ssh %s: aborted", mc.Host)
return
}
}
func (m *Monitor) setSSHError(st *machineState, msg string) {
st.mu.Lock()
defer st.mu.Unlock()
st.sshResult = &sshResultState{
Ok: false,
Error: msg,
LastRun: time.Now(),
ExitCode: nil,
}
}
type parsedMetric struct {
name string
value float64
unit string
}
func parseMetrics(out string) []parsedMetric {
metrics := []parsedMetric{}
for _, line := range strings.Split(out, "\n") {
if strings.Contains(line, "---") {
break
}
parts := strings.SplitN(line, ":", 3)
if len(parts) < 3 {
continue
}
name := strings.TrimSpace(parts[0])
status := strings.TrimSpace(parts[1])
if name == "" || status != "metric" {
continue
}
val := strings.TrimSpace(parts[2])
num := 0.0
unit := ""
for i, r := range val {
if (r >= '0' && r <= '9') || r == '.' || r == '-' || r == '+' {
continue
}
if p, err := strconv.ParseFloat(strings.TrimSpace(val[:i]), 64); err == nil {
num = p
unit = strings.TrimSpace(val[i:])
}
break
}
if unit == "" && val != "" {
if p, err := strconv.ParseFloat(val, 64); err == nil {
num = p
} else {
continue
}
}
metrics = append(metrics, parsedMetric{name: name, value: num, unit: unit})
}
return metrics
}
func (m *Monitor) accumulateMetrics(st *machineState, machineID, out string, now time.Time) {
metrics := parseMetrics(out)
if len(metrics) == 0 {
return
}
ts := now.Unix()
st.mu.Lock()
if st.metrics == nil {
st.metrics = make(map[string]metricState, len(metrics))
}
for _, p := range metrics {
st.metrics[p.name] = metricState{value: p.value, unit: p.unit, updated: now}
}
st.mu.Unlock()
if m.history != nil {
for _, p := range metrics {
m.history.RecordMetric(machineID, p.name, ts, p.value, p.unit)
}
}
}
func (m *Monitor) StatusPayload() statusPayload {
var dbSize, dbSizeMonth int64
var dbRows int
if m.history != nil {
dbSize = m.history.Size()
rows, firstTS := m.history.Stats()
dbRows = rows
if rows > 0 && firstTS > 0 {
elapsed := time.Now().Unix() - firstTS
if elapsed < 1 {
elapsed = 1
}
perSec := float64(dbSize) / float64(elapsed)
dbSizeMonth = dbSize + int64(perSec*float64(30*24*3600))
}
}
p := statusPayload{
Title: m.cfg.Title,
Interactive: m.interactive,
SharedMetricWindow: m.cfg.SharedMetricWindow,
Stats: statsPayload{
PingInterval: m.cfg.Ping.Interval.Duration.String(),
PingTimeout: m.cfg.Ping.Timeout.Duration.String(),
TCPPort: m.cfg.Ping.TCPPort,
SSHInterval: m.cfg.SSH.Interval.Duration.String(),
SSHEnabled: m.cfg.SSH.Key != "" && m.cfg.SSH.Script != "",
DBSize: dbSize,
DBRows: dbRows,
DBSizeMonth: dbSizeMonth,
},
Groups: []groupStatus{},
Machines: []machineStatus{},
}
if len(m.cfg.Metrics) > 0 {
p.MetricRanges = make(map[string]metricRange, len(m.cfg.Metrics))
for name, r := range m.cfg.Metrics {
p.MetricRanges[name] = metricRange{Min: r.Min, Max: r.Max}
}
}
grouped := map[string][]machineStatus{}
order := []string{}
appendStatus := func(mc Machine, group string) {
st := m.state(mc)
ms := machineStatus{
ID: mc.ID,
Name: mc.Name,
Host: mc.Host,
SSHConfigured: mc.SSH.Key != "" && mc.SSH.Script != "",
}
if st != nil {
st.mu.RLock()
ms.Status = st.status
ms.ICMP = st.icmp
ms.TCP = st.tcp
ms.IP = st.ip
ms.IPs = append([]string{}, st.ips...)
if !st.lastPing.IsZero() {
t := st.lastPing
ms.LastPing = &t
}
if st.sshResult != nil {
s := *st.sshResult
ms.SSH = &s
}
if len(st.metrics) > 0 {
ms.Metrics = make([]metricStatus, 0, len(st.metrics))
for name, mtr := range st.metrics {
ms.Metrics = append(ms.Metrics, metricStatus{
Name: name,
Value: mtr.value,
Unit: mtr.unit,
Updated: &mtr.updated,
})
}
}
st.mu.RUnlock()
} else {
ms.Status = statusUnknown
ms.ICMP = checkNA
ms.TCP = checkNA
}
if group == "" {
p.Machines = append(p.Machines, ms)
return
}
if _, ok := grouped[group]; !ok {
order = append(order, group)
}
grouped[group] = append(grouped[group], ms)
}
for _, g := range m.cfg.Groups {
for _, mc := range g.Machines {
appendStatus(resolveMachine(mc, g.Name, m.cfg), g.Name)
}
}
for _, mc := range m.cfg.Machines {
if mc.Group == "" {
appendStatus(mc, "")
}
}
for _, name := range order {
p.Groups = append(p.Groups, groupStatus{Name: name, Machines: grouped[name]})
}
return p
}
func (m *Monitor) statusHandler(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(m.StatusPayload())
}
func (m *Monitor) historyHandler(w http.ResponseWriter, r *http.Request) {
machine := r.PathValue("machine")
if machine == "" {
http.Error(w, "missing machine", http.StatusBadRequest)
return
}
limit, err := strconv.Atoi(r.URL.Query().Get("limit"))
if err != nil {
limit = 500
}
entries, err := m.history.Query(machine, limit)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(entries)
}
func (m *Monitor) metricsHandler(w http.ResponseWriter, r *http.Request) {
if m.history == nil {
http.Error(w, "no history db", http.StatusInternalServerError)
return
}
machine := r.PathValue("machine")
if machine == "" {
http.Error(w, "missing machine", http.StatusBadRequest)
return
}
q := r.URL.Query()
maxPoints := 100
if v := q.Get("max_points"); v != "" {
if n, err := strconv.Atoi(v); err == nil {
maxPoints = n
}
} else if v := q.Get("limit"); v != "" {
if n, err := strconv.Atoi(v); err == nil {
maxPoints = n
}
}
var minTS, maxTS int64
if v := q.Get("min"); v != "" {
minTS, _ = strconv.ParseInt(v, 10, 64)
}
if v := q.Get("max"); v != "" {
maxTS, _ = strconv.ParseInt(v, 10, 64)
}
entries, err := m.history.QueryMetrics(machine, q.Get("name"), minTS, maxTS, maxPoints, q.Get("shared") != "")
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(entries)
}
func (m *Monitor) metricsAggregateHandler(w http.ResponseWriter, r *http.Request) {
if m.history == nil {
http.Error(w, "no history db", http.StatusInternalServerError)
return
}
q := r.URL.Query()
var minTS, maxTS int64
if v := q.Get("min"); v != "" {
minTS, _ = strconv.ParseInt(v, 10, 64)
}
if v := q.Get("max"); v != "" {
maxTS, _ = strconv.ParseInt(v, 10, 64)
}
metrics, window, err := m.history.QueryMetricsAll(minTS, maxTS)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(struct {
Window [2]int64 `json:"window"`
Metrics []MetricAggregate `json:"metrics"`
}{window, metrics})
}
func (m *Monitor) refreshHandler(w http.ResponseWriter, r *http.Request) {
if !m.interactive {
http.Error(w, "interactive mode disabled", http.StatusForbidden)
return
}
id := r.PathValue("machine")
mc, ok := m.machinesByID[id]
if !ok {
http.Error(w, "machine not found", http.StatusNotFound)
return
}
log.Printf("interactive refresh scheduled for %s (%s)", mc.Name, mc.Host)
m.Refresh(mc)
w.WriteHeader(http.StatusAccepted)
}