Files
pwa-homelab-mon/server/main.go
T
alexey.bagno 8eae79af06
Deploy / deploy (push) Successful in 12s
fix parsing systemd
2026-08-31 16:48:16 +05:00

612 lines
19 KiB
Go

package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"fmt"
"log"
"net/http"
"os"
"os/exec"
"path/filepath"
"regexp"
"sort"
"strings"
"sync"
"time"
"connectrpc.com/connect"
"github.com/labstack/echo/v4"
"github.com/labstack/echo/v4/middleware"
"gopkg.in/yaml.v3"
homelabv1 "pwa-homelab-mon/gen/homelab/v1"
"pwa-homelab-mon/gen/homelab/v1/homelabv1connect"
)
// --- config ---
type Item struct {
Name string `yaml:"name"`
Script string `yaml:"script"`
Refresh int `yaml:"refresh"` // monitor: seconds
Data string `yaml:"data"` // monitor: text | percentage
Args int `yaml:"args"` // exec: arg count
ArgTypes []string `yaml:"arg_types"` // exec: text | number per arg
}
type Link struct {
Name string `yaml:"name"`
URL string `yaml:"url"`
}
type Config struct {
Links []Link `yaml:"links"`
Monitor []Item `yaml:"monitor"`
Exec []Item `yaml:"exec"`
}
func loadConfig(path string) (*Config, error) {
f, err := os.Open(path)
if err != nil {
return nil, err
}
defer f.Close()
var cfg Config
dec := yaml.NewDecoder(f)
dec.KnownFields(true)
if err := dec.Decode(&cfg); err != nil {
return nil, fmt.Errorf("parse %s: %w", path, err)
}
seen := map[string]bool{} // item names are global keys (monVals, RunScript)
for _, l := range cfg.Links {
if seen[l.Name] {
return nil, fmt.Errorf("duplicate name %q", l.Name)
}
seen[l.Name] = true
}
for _, list := range [][]Item{cfg.Monitor, cfg.Exec} {
for _, it := range list {
if seen[it.Name] {
return nil, fmt.Errorf("duplicate name %q", it.Name)
}
seen[it.Name] = true
}
}
for _, it := range cfg.Monitor {
if it.Refresh <= 0 {
return nil, fmt.Errorf("monitor %q: refresh must be > 0", it.Name)
}
if it.Data != "text" && it.Data != "percentage" {
return nil, fmt.Errorf("monitor %q: data must be text|percentage", it.Name)
}
}
for _, it := range cfg.Exec {
if len(it.ArgTypes) != it.Args {
return nil, fmt.Errorf("exec %q: arg_types must have args entries", it.Name)
}
for _, t := range it.ArgTypes {
if t != "text" && t != "number" {
return nil, fmt.Errorf("exec %q: arg type must be text|number", it.Name)
}
}
}
return &cfg, nil
}
// --- service ---
type server struct {
cfg *Config
exec map[string]*Item // by item name, read-only after start
monVals sync.Map // monitor item name -> string (last output)
docsDir string
srvDir string // dir with *.service units (~/srv on the box)
probeErr sync.Map // probe key -> last logged error; poll errors log once per change, not every 3s tick
}
const execTimeout = 10 * time.Second
func runScript(s *Item, args ...string) (string, error) {
ctx, cancel := context.WithTimeout(context.Background(), execTimeout)
defer cancel()
out, err := exec.CommandContext(ctx, s.Script, args...).CombinedOutput()
if err != nil {
return "", fmt.Errorf("%v: %s", err, out)
}
return string(out), nil
}
func (s *server) startMonitors() {
for i := range s.cfg.Monitor {
sc := &s.cfg.Monitor[i]
log.Printf("monitor started: %s (%s) every %ds", sc.Name, sc.Script, sc.Refresh)
go func(sc *Item) {
for {
v, err := runScript(sc)
if err != nil {
log.Printf("monitor %s: %v", sc.Name, err)
v = "error: " + err.Error()
} else {
v = strings.TrimSpace(v)
log.Printf("monitor %s -> %s", sc.Name, v)
}
s.monVals.Store(sc.Name, v)
time.Sleep(time.Duration(sc.Refresh) * time.Second)
}
}(sc)
}
}
// --- docs ---
func docHash(b []byte) string {
h := sha256.Sum256(b)
return hex.EncodeToString(h[:])
}
// docPath resolves a doc name safely: plain file name only, no separators, always inside docsDir.
func (s *server) docPath(name string) (string, error) {
if filepath.Base(name) != name || filepath.Ext(name) != ".md" {
return "", fmt.Errorf("bad doc name %q", name)
}
p := filepath.Join(s.docsDir, name)
if filepath.Dir(p) != s.docsDir {
return "", fmt.Errorf("bad doc name %q", name)
}
return p, nil
}
func (s *server) ListDocs(context.Context, *connect.Request[homelabv1.ListDocsRequest]) (*connect.Response[homelabv1.ListDocsResponse], error) {
resp := &homelabv1.ListDocsResponse{}
entries, err := os.ReadDir(s.docsDir)
if err != nil {
return connect.NewResponse(resp), nil // no docs dir = no docs
}
for _, e := range entries {
if e.IsDir() || !strings.HasSuffix(e.Name(), ".md") {
continue
}
b, err := os.ReadFile(filepath.Join(s.docsDir, e.Name()))
if err != nil {
continue
}
resp.Docs = append(resp.Docs, &homelabv1.DocMeta{Name: e.Name(), Hash: docHash(b)})
}
return connect.NewResponse(resp), nil
}
func (s *server) GetDoc(_ context.Context, req *connect.Request[homelabv1.GetDocRequest]) (*connect.Response[homelabv1.GetDocResponse], error) {
p, err := s.docPath(req.Msg.GetName())
if err != nil {
return nil, connect.NewError(connect.CodeInvalidArgument, err)
}
b, err := os.ReadFile(p)
if err != nil {
return nil, connect.NewError(connect.CodeNotFound, err)
}
return connect.NewResponse(&homelabv1.GetDocResponse{Content: string(b)}), nil
}
func (s *server) SaveDoc(_ context.Context, req *connect.Request[homelabv1.SaveDocRequest]) (*connect.Response[homelabv1.SaveDocResponse], error) {
p, err := s.docPath(req.Msg.GetName())
if err != nil {
return nil, connect.NewError(connect.CodeInvalidArgument, err)
}
content := []byte(req.Msg.GetContent())
// ponytail: client authority — last write wins, no base-hash check; add optimistic locking if concurrent edits appear
if err := os.WriteFile(p, content, 0o644); err != nil {
return nil, connect.NewError(connect.CodeInternal, err)
}
log.Printf("doc saved: %s (%d bytes)", filepath.Base(p), len(content))
return connect.NewResponse(&homelabv1.SaveDocResponse{Hash: docHash(content)}), nil
}
// --- services: systemd user units + docker containers, one flat list ---
// Service names double as file/unit/container selectors — no separators, no traversal.
var svcNameRe = regexp.MustCompile(`^[a-z0-9_.-]+$`)
const svcTimeout = 8 * time.Second
func (s *server) services(ctx context.Context) []*homelabv1.Service {
return append(s.systemdServices(ctx), s.dockerServices(ctx)...)
}
// logProbeChange logs poll-path probe failures once per change: the poll repeats every 3s per client,
// and the same error every tick is noise. msg "" = probe healthy, logs a single recovery line.
func (s *server) logProbeChange(key, msg string) {
prev, loaded := s.probeErr.LoadOrStore(key, msg)
if !loaded {
if msg != "" {
log.Printf("probe %s: %s", key, msg)
}
return
}
if prev.(string) == msg {
return
}
s.probeErr.Store(key, msg)
if msg == "" {
log.Printf("probe %s: recovered", key)
} else {
log.Printf("probe %s: %s", key, msg)
}
}
func (s *server) units() ([]string, error) {
entries, err := os.ReadDir(s.srvDir)
if err != nil {
return nil, err
}
units := make([]string, 0, len(entries))
for _, e := range entries {
if n := strings.TrimSuffix(e.Name(), ".service"); n != e.Name() && svcNameRe.MatchString(n) {
units = append(units, n)
}
}
sort.Strings(units)
return units, nil
}
func (s *server) systemdServices(ctx context.Context) []*homelabv1.Service {
units, err := s.units()
if err != nil {
s.logProbeChange("units", err.Error())
return []*homelabv1.Service{{Kind: homelabv1.ServiceKind_SYSTEMD, Name: s.srvDir, Error: err.Error()}}
}
s.logProbeChange("units", "")
svcs := make([]*homelabv1.Service, 0, len(units))
for _, u := range units {
svcs = append(svcs, s.systemdStatus(ctx, u))
}
return svcs
}
// ponytail: one systemctl per unit (~26 per 3s poll, ~10ms each). Batch into a single
// `systemctl show u1 u2 ... -p ...` if the poll ever shows up in a slow trace.
func (s *server) systemdStatus(ctx context.Context, unit string) *homelabv1.Service {
svc := &homelabv1.Service{Kind: homelabv1.ServiceKind_SYSTEMD, Name: unit}
ctx, cancel := context.WithTimeout(ctx, svcTimeout)
defer cancel()
// show is the machine-readable probe (always exit 0, unlike status). LoadState=not-found means the user
// disabled the unit on purpose — a state, not an error.
out, err := exec.CommandContext(ctx, "systemctl", "--user", "show", unit+".service",
"-p", "ActiveState,SubState,LoadState,UnitFileState").Output()
// systemctl prints properties in its own fixed table order (not the -p order), so parse by key
props := map[string]string{}
for _, line := range strings.Split(string(out), "\n") {
if k, v, ok := strings.Cut(strings.TrimSpace(line), "="); ok {
props[k] = v
}
}
svc.State, svc.Substate = props["ActiveState"], props["SubState"]
svc.LoadState, svc.UnitFileState = props["LoadState"], props["UnitFileState"]
if svc.LoadState == "" {
svc.LoadState = "not-found"
}
if err != nil {
svc.Error = err.Error()
s.logProbeChange("systemctl/"+unit, svc.Error)
} else {
s.logProbeChange("systemctl/"+unit, "")
}
return svc
}
// dockerC is one container as seen by `docker inspect`.
type dockerC struct {
name, project, status, policy string
}
func dockerOut(ctx context.Context, args ...string) (string, error) {
ctx, cancel := context.WithTimeout(ctx, svcTimeout)
defer cancel()
out, err := exec.CommandContext(ctx, "docker", args...).Output()
return string(out), err
}
func (s *server) dockerContainers(ctx context.Context) ([]dockerC, error) {
ids, err := dockerOut(ctx, "ps", "-aq")
if err != nil {
return nil, err
}
if strings.TrimSpace(ids) == "" {
return nil, nil
}
out, err := dockerOut(ctx, append([]string{"inspect", "-f",
`{{index .Config.Labels "com.docker.compose.project"}}|{{.Name}}|{{.State.Status}}|{{.HostConfig.RestartPolicy.Name}}`},
strings.Fields(ids)...)...)
if err != nil {
return nil, err
}
var cs []dockerC
for _, line := range strings.Split(strings.TrimSpace(out), "\n") {
f := strings.Split(line, "|")
if len(f) != 4 {
continue
}
cs = append(cs, dockerC{project: f[0], name: strings.TrimPrefix(f[1], "/"), status: f[2], policy: f[3]})
}
sort.Slice(cs, func(i, j int) bool { return cs[i].name < cs[j].name })
return cs, nil
}
// working drops one-shot helpers (restart policy no + already exited: *-init, migrations). They belong to
// their project's card but must not drag its aggregate down, and must not be re-run by the start button.
func working(cs []dockerC) []dockerC {
out := make([]dockerC, 0, len(cs))
for _, c := range cs {
if c.policy == "no" && c.status != "running" {
continue
}
out = append(out, c)
}
return out
}
func (s *server) dockerServices(ctx context.Context) []*homelabv1.Service {
cs, err := s.dockerContainers(ctx)
if err != nil {
s.logProbeChange("docker", err.Error())
return []*homelabv1.Service{{Kind: homelabv1.ServiceKind_DOCKER, Name: "docker", Error: err.Error()}}
}
s.logProbeChange("docker", "")
groups := map[string][]dockerC{}
var keys []string
for _, c := range cs {
key := c.project
if key == "" {
key = c.name // docker run without a compose label -> its own card
}
if _, ok := groups[key]; !ok {
keys = append(keys, key)
}
groups[key] = append(groups[key], c)
}
sort.Strings(keys)
svcs := make([]*homelabv1.Service, 0, len(keys))
for _, k := range keys {
if svc := dockerCard(k, working(groups[k])); svc != nil {
svcs = append(svcs, svc)
}
}
return svcs
}
// dockerCard aggregates a project: all running -> active, none -> inactive, else partial (up/total).
func dockerCard(name string, cs []dockerC) *homelabv1.Service {
if len(cs) == 0 {
return nil
}
up := 0
names := make([]string, 0, len(cs))
for _, c := range cs {
if c.status == "running" {
up++
}
names = append(names, c.name)
}
state := "partial"
switch {
case up == len(cs):
state = "active"
case up == 0:
state = "inactive"
}
return &homelabv1.Service{Kind: homelabv1.ServiceKind_DOCKER, Name: name, State: state,
Substate: fmt.Sprintf("%d/%d", up, len(cs)), Containers: names}
}
// --- content RPC ---
func (s *server) ListContent(ctx context.Context, _ *connect.Request[homelabv1.ListContentRequest]) (*connect.Response[homelabv1.ListContentResponse], error) {
resp := &homelabv1.ListContentResponse{}
for _, l := range s.cfg.Links {
resp.Links = append(resp.Links, &homelabv1.Link{Name: l.Name, Url: l.URL})
}
for i := range s.cfg.Monitor {
it := &s.cfg.Monitor[i]
val, _ := s.monVals.Load(it.Name)
v, _ := val.(string)
resp.Monitors = append(resp.Monitors, &homelabv1.Item{Name: it.Name, Data: it.Data, Value: v})
}
for i := range s.cfg.Exec {
it := &s.cfg.Exec[i]
resp.Exec = append(resp.Exec, &homelabv1.Item{Name: it.Name, Args: int32(it.Args), ArgTypes: it.ArgTypes})
}
resp.Services = s.services(ctx)
return connect.NewResponse(resp), nil
}
// ServiceAction runs a fixed systemctl/docker command; only `name` is user input and it is regex-checked,
// never interpolated into a shell string.
func (s *server) ServiceAction(ctx context.Context, req *connect.Request[homelabv1.ServiceActionRequest]) (*connect.Response[homelabv1.ServiceActionResponse], error) {
name, kind := req.Msg.GetName(), req.Msg.GetKind()
if !svcNameRe.MatchString(name) {
return nil, connect.NewError(connect.CodeInvalidArgument, fmt.Errorf("bad service name %q", name))
}
var args []string
switch req.Msg.GetOp() {
case homelabv1.ServiceOp_START:
args = []string{"start"}
case homelabv1.ServiceOp_STOP:
args = []string{"stop"}
case homelabv1.ServiceOp_RESTART:
args = []string{"restart"}
default:
return nil, connect.NewError(connect.CodeInvalidArgument, fmt.Errorf("bad op for %q", name))
}
ctx, cancel := context.WithTimeout(ctx, svcTimeout)
defer cancel()
var cmd *exec.Cmd
if kind == homelabv1.ServiceKind_DOCKER {
targets, err := s.dockerTargets(ctx, name, args[0])
if err != nil {
return nil, connect.NewError(connect.CodeInternal, err)
}
if len(targets) == 0 {
return connect.NewResponse(&homelabv1.ServiceActionResponse{Output: "no containers to " + args[0]}), nil
}
cmd = exec.CommandContext(ctx, "docker", append([]string{args[0]}, targets...)...)
} else {
cmd = exec.CommandContext(ctx, "systemctl", "--user", args[0], name+".service")
}
out, err := cmd.CombinedOutput()
if err != nil {
log.Printf("service %s/%s %s: ERROR %v: %s", kind, name, args[0], err, out)
return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("%v: %s", err, out))
}
log.Printf("service %s/%s %s -> %s", kind, name, args[0], strings.TrimSpace(string(out)))
return connect.NewResponse(&homelabv1.ServiceActionResponse{Output: strings.TrimSpace(string(out))}), nil
}
// dockerTargets resolves a card name (compose project, or bare container name) to the containers an op applies to.
func (s *server) dockerTargets(ctx context.Context, name, op string) ([]string, error) {
cs, err := s.dockerContainers(ctx)
if err != nil {
return nil, err
}
var targets []string
for _, c := range working(cs) {
if c.project != name && c.name != name {
continue
}
if running := c.status == "running"; (op == "start" && !running) || (op != "start" && running) {
targets = append(targets, c.name)
}
}
return targets, nil
}
// ServiceInfo returns raw status output for the modal: `systemctl status` / `docker ps -a` for the card.
func (s *server) ServiceInfo(ctx context.Context, req *connect.Request[homelabv1.ServiceInfoRequest]) (*connect.Response[homelabv1.ServiceInfoResponse], error) {
name, kind := req.Msg.GetName(), req.Msg.GetKind()
if !svcNameRe.MatchString(name) {
return nil, connect.NewError(connect.CodeInvalidArgument, fmt.Errorf("bad service name %q", name))
}
ctx, cancel := context.WithTimeout(ctx, svcTimeout)
defer cancel()
var out string
var err error
if kind == homelabv1.ServiceKind_DOCKER {
// compose project first, bare container (no label) as fallback
out, err = dockerOut(ctx, "ps", "-a", "--filter", "label=com.docker.compose.project="+name)
if strings.TrimSpace(out) == "" {
out, err = dockerOut(ctx, "ps", "-a", "--filter", "name="+name)
}
} else {
// status exits non-zero for inactive units (code 3) — that is valid output, not an error
c := exec.CommandContext(ctx, "systemctl", "--user", "status", name+".service", "--no-pager", "-n", "0")
b, cerr := c.CombinedOutput()
out, err = string(b), cerr
}
if err != nil && strings.TrimSpace(out) == "" {
log.Printf("service info %s/%s: ERROR %v", kind, name, err)
return nil, connect.NewError(connect.CodeInternal, err)
}
log.Printf("service info %s/%s ok", kind, name)
return connect.NewResponse(&homelabv1.ServiceInfoResponse{Output: strings.TrimRight(out, "\n")}), nil
}
func (s *server) RunScript(_ context.Context, req *connect.Request[homelabv1.RunScriptRequest]) (*connect.Response[homelabv1.RunScriptResponse], error) {
it, ok := s.exec[req.Msg.GetName()]
if !ok {
return nil, connect.NewError(connect.CodeNotFound, fmt.Errorf("no exec item %q", req.Msg.GetName()))
}
if len(req.Msg.GetArgs()) != it.Args {
return nil, connect.NewError(connect.CodeInvalidArgument, fmt.Errorf("%q takes %d args, got %d", it.Name, it.Args, len(req.Msg.GetArgs())))
}
out, err := runScript(it, req.Msg.GetArgs()...)
if err != nil {
log.Printf("script %s args=%v: ERROR %v", it.Name, req.Msg.GetArgs(), err)
return nil, connect.NewError(connect.CodeInternal, err)
}
log.Printf("script %s args=%v -> %s", it.Name, req.Msg.GetArgs(), strings.TrimSpace(out))
return connect.NewResponse(&homelabv1.RunScriptResponse{Output: out}), nil
}
// --- main ---
func main() {
cfgPath := os.Getenv("CONFIG")
if cfgPath == "" {
cfgPath = "config.yaml"
}
cfg, err := loadConfig(cfgPath)
if err != nil {
fmt.Fprintln(os.Stderr, "config:", err)
os.Exit(1)
}
docsDir := "docs"
if err := os.MkdirAll(docsDir, 0o755); err != nil {
fmt.Fprintln(os.Stderr, "docs:", err)
os.Exit(1)
}
svc := &server{cfg: cfg, exec: map[string]*Item{}, docsDir: docsDir, srvDir: srvDir()}
for i := range cfg.Exec {
svc.exec[cfg.Exec[i].Name] = &cfg.Exec[i]
}
svc.startMonitors()
log.Printf("config %s: %d links, %d monitors, %d exec; units from %s", cfgPath, len(cfg.Links), len(cfg.Monitor), len(cfg.Exec), svc.srvDir)
e := echo.New()
e.HideBanner = true
// HTTP request log, except the ListContent heartbeat (every 3s per client — pure noise).
e.Use(middleware.LoggerWithConfig(middleware.LoggerConfig{
Skipper: func(c echo.Context) bool {
return c.Request().URL.Path == "/homelab.v1.HomelabService/ListContent"
},
}))
// ConnectRPC handler on the mux (matches /homelab.v1.HomelabService/<Method>).
mux := http.NewServeMux()
mux.Handle(homelabv1connect.NewHomelabServiceHandler(svc))
e.Any("/homelab.v1.HomelabService/*", echo.WrapHandler(mux))
// sw.js and index.html must never be cached: stale index = stale asset references = iOS never updates.
e.Use(func(next echo.HandlerFunc) echo.HandlerFunc {
return func(c echo.Context) error {
p := c.Request().URL.Path
if p == "/sw.js" || p == "/manifest.webmanifest" || p == "/index.html" || p == "/" {
c.Response().Header().Set("Cache-Control", "no-cache")
}
return next(c)
}
})
// Static PWA from dist/ (built by make web); HTML5 mode falls back to index.html for SPA routes.
e.Use(middleware.StaticWithConfig(middleware.StaticConfig{
Root: "dist",
HTML5: true,
}))
e.Logger.Fatal(e.Start(port()))
}
func port() string {
if p := os.Getenv("PORT"); p != "" {
return ":" + p
}
return ":8080"
}
// srvDir: units live in ~/srv on the box, which is outside WorkingDirectory, so it is not repo-relative
// (unlike dist/ and docs/). SRV_DIR overrides it for dev/tests.
func srvDir() string {
if d := os.Getenv("SRV_DIR"); d != "" {
return d
}
home, err := os.UserHomeDir()
if err != nil {
return "srv"
}
return filepath.Join(home, "srv")
}