#!/usr/bin/env python3 """Pulse agent: a read-only health snapshot, sent once a minute. It reads /proc, df, ps, docker, systemctl and journalctl. It never changes anything on the server and never sends command lines, environment variables or file contents. Error-log lines are redacted before they leave the host. Config: /etc/pulse/config (PULSE_URL, PULSE_TOKEN), optional /etc/pulse/domains (one hostname per line, for TLS certificate expiry checks). """ import json import os import re import shutil import socket import ssl import subprocess import time import urllib.request VERSION = "0.1.0" CONFIG = "/etc/pulse/config" DOMAINS = "/etc/pulse/domains" STATE = "/var/lib/pulse/state.json" REDACT = [ (re.compile(r"(?i)\b(authorization|proxy-authorization|cookie|set-cookie)\b\s*[:=].*"), r"\1: [redacted]"), (re.compile(r"(?i)\b(bearer|basic)\s+[A-Za-z0-9._~+/=-]+"), r"\1 [redacted]"), (re.compile(r"(?i)(password|passwd|pwd|secret|token|api[_-]?key)(\s*[=:]\s*|\s+)\S+"), r"\1\2[redacted]"), (re.compile(r"[A-Za-z0-9+/_=-]{32,}"), "[redacted]"), (re.compile(r"\b[\w.+-]+@[\w-]+\.[\w.-]+\b"), "[email]"), ] def sh(*cmd, timeout=10): try: out = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout) return out.stdout if out.returncode == 0 else "" except (OSError, subprocess.TimeoutExpired): return "" def read(path): try: with open(path) as f: return f.read() except OSError: return "" def cpu_percent(): def sample(): parts = read("/proc/stat").split("\n", 1)[0].split()[1:] values = [int(x) for x in parts] idle = values[3] + (values[4] if len(values) > 4 else 0) return idle, sum(values) i1, t1 = sample() time.sleep(1) i2, t2 = sample() return round(100 * (1 - (i2 - i1) / max(t2 - t1, 1)), 1) def memory(): info = {} for line in read("/proc/meminfo").splitlines(): key, _, rest = line.partition(":") info[key] = int(rest.split()[0]) * 1024 if rest.split() else 0 total, avail = info.get("MemTotal", 0), info.get("MemAvailable", 0) return {"total": total, "available": avail, "used_pct": round(100 * (1 - avail / total), 1) if total else None, "swap_total": info.get("SwapTotal", 0), "swap_free": info.get("SwapFree", 0)} def disks(): result = [] for line in sh("df", "-P", "-k", "-x", "tmpfs", "-x", "devtmpfs", "-x", "overlay", "-x", "squashfs").splitlines()[1:]: f = line.split() if len(f) >= 6 and f[0].startswith("/"): result.append({"mount": f[5], "size": int(f[1]) * 1024, "used": int(f[2]) * 1024, "used_pct": int(f[4].rstrip("%"))}) return result def processes(): rows = [] for line in sh("ps", "-eo", "comm,%cpu,%mem", "--sort=-%cpu", "--no-headers").splitlines()[:8]: parts = line.split() if len(parts) >= 3: rows.append({"name": " ".join(parts[:-2])[:40], "cpu": float(parts[-2]), "mem": float(parts[-1])}) return rows def containers(): if not shutil.which("docker"): return None rows = [] for line in sh("docker", "ps", "-a", "--format", "{{.Names}}|{{.Status}}|{{.Image}}").splitlines(): name, _, rest = line.partition("|") status, _, image = rest.partition("|") rows.append({"name": name, "status": status, "image": image.split("@")[0][:80], "running": status.startswith("Up"), "restarting": status.startswith("Restarting")}) return rows def docker_disk(): if not shutil.which("docker"): return None rows = [] for line in sh("docker", "system", "df", "--format", "{{.Type}}|{{.Size}}|{{.Reclaimable}}", timeout=30).splitlines(): kind, size, reclaim = (line.split("|") + ["", "", ""])[:3] rows.append({"type": kind, "size": size, "reclaimable": reclaim}) return rows def failed_units(): out = sh("systemctl", "--failed", "--no-legend", "--plain") return [line.split()[0] for line in out.splitlines() if line.strip()] def recent_errors(): out = sh("journalctl", "-p", "err", "--since", "-15min", "-o", "short-iso", "--no-pager", "-n", "25") lines = [] for line in out.splitlines(): if line.startswith("--"): continue for pattern, repl in REDACT: line = pattern.sub(repl, line) lines.append(line[:300]) return lines def cert_expiry(host): try: ctx = ssl.create_default_context() with socket.create_connection((host, 443), timeout=5) as s: with ctx.wrap_socket(s, server_hostname=host) as t: expires = ssl.cert_time_to_seconds(t.getpeercert()["notAfter"]) return {"host": host, "days_left": round((expires - time.time()) / 86400, 1), "ok": True} except Exception as e: # report the failure kind, not the details return {"host": host, "days_left": None, "ok": False, "error": type(e).__name__} def load_state(): try: return json.loads(read(STATE) or "{}") except ValueError: return {} def save_state(state): os.makedirs(os.path.dirname(STATE), exist_ok=True) with open(STATE, "w") as f: json.dump(state, f) def config(): values = {} for line in read(CONFIG).splitlines(): key, sep, value = line.partition("=") if sep: values[key.strip()] = value.strip().strip('"') return values def main(): cfg = config() url, token = cfg.get("PULSE_URL"), cfg.get("PULSE_TOKEN") if not url or not token: raise SystemExit("pulse-agent: missing PULSE_URL or PULSE_TOKEN in " + CONFIG) state, now = load_state(), time.time() os_name = next((l.split("=", 1)[1].strip('"') for l in read("/etc/os-release").splitlines() if l.startswith("PRETTY_NAME=")), "Linux") load = read("/proc/loadavg").split() snapshot = { "agent": VERSION, "ts": now, "hostname": socket.gethostname(), "os": os_name, "kernel": os.uname().release, "uptime": float(read("/proc/uptime").split()[0] or 0), "cpu": {"cores": os.cpu_count(), "used_pct": cpu_percent(), "load": [float(x) for x in load[:3]] if load else None}, "memory": memory(), "disks": disks(), "processes": processes(), "containers": containers(), "failed_units": failed_units(), "errors": recent_errors(), } # Slower checks run less often and are re-sent from the cached state in between. if now - state.get("docker_disk_ts", 0) > 900: state["docker_disk"], state["docker_disk_ts"] = docker_disk(), now if now - state.get("certs_ts", 0) > 3600: hosts = [h.strip() for h in read(DOMAINS).splitlines() if h.strip() and not h.startswith("#")][:40] state["certs"], state["certs_ts"] = [cert_expiry(h) for h in hosts], now snapshot["docker_disk"], snapshot["certs"] = state.get("docker_disk"), state.get("certs", []) save_state(state) request = urllib.request.Request(url.rstrip("/") + "/api/ingest", data=json.dumps(snapshot).encode(), headers={"Authorization": "Bearer " + token, "Content-Type": "application/json", "User-Agent": "pulse-agent/" + VERSION}, method="POST") with urllib.request.urlopen(request, timeout=20) as response: response.read() if __name__ == "__main__": main()