743 lines
29 KiB
Python
Executable File
743 lines
29 KiB
Python
Executable File
#!/usr/bin/env python3
|
|
"""
|
|
semprini-maintainer — autonomous stack maintenance daemon.
|
|
|
|
Runs as a systemd service. Periodically:
|
|
* (frequent) reads Uptime Kuma to learn which services are tracked and their
|
|
up/down status, and — when something is down — invokes the
|
|
canonical stack-support agent (Claude Code, headless) to
|
|
diagnose and remediate fully autonomously.
|
|
* (infrequent) resolves the containers behind the tracked services, finds the
|
|
compose file that owns each, and invokes the stack-support
|
|
agent to research + apply upgrades (test, then rollback on
|
|
failure). Ambiguous / breaking upgrades are escalated to the
|
|
admin over Matrix and resumed once they reply.
|
|
|
|
All operator communication happens through the core stack's Matrix
|
|
(chat.semprini.me) via a dedicated bot account.
|
|
|
|
The scope ("which services to maintain") is whatever Uptime Kuma is tracking at
|
|
the moment the health check runs — it is re-read every cycle, never hard-coded.
|
|
|
|
Design notes
|
|
------------
|
|
* Pure Python stdlib only (urllib, sqlite via a throwaway container, subprocess)
|
|
so the service deploys without a venv.
|
|
* The daemon is intentionally "dumb": it gathers state and hands the fuzzy work
|
|
(mapping a monitor to a container, choosing a fix, planning an upgrade) to the
|
|
LLM agent defined in agents/stack-support.agent.md. The agent returns a small
|
|
JSON envelope which the daemon relays / acts on.
|
|
"""
|
|
|
|
import argparse
|
|
import json
|
|
import os
|
|
import re
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
import traceback
|
|
import urllib.parse
|
|
import urllib.request
|
|
from datetime import datetime, timezone
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Config / state
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
DEFAULTS = {
|
|
"repo_dir": "/home/paul/Dev/semprini-core",
|
|
# Extra directories the agent may read/operate on (other projects that use
|
|
# the core stack, e.g. semprini-blog). Passed to claude via --add-dir.
|
|
"extra_dirs": ["/home/paul/Dev"],
|
|
"agent_file": "agents/stack-support.agent.md",
|
|
"runbook_file": "docs/stack-support-runbook.md",
|
|
"claude_bin": "/home/paul/.local/bin/claude",
|
|
# Pin a standard 200K-window model. Leaving this None uses the account
|
|
# default, which may be a 1M-context ([1m]) model; once a cycle's context
|
|
# crosses 200K that escalates to the 1M tier and fails on subscription
|
|
# plans with "Usage credits are required for long context requests" (429).
|
|
"claude_model": "claude-sonnet-4-6",
|
|
"claude_extra_args": [],
|
|
"agent_timeout_seconds": 1800,
|
|
"kuma_volume": "uptime-kuma-data",
|
|
"sqlite_image": "keinos/sqlite3",
|
|
# autonomy: "full" (fix + upgrade), "fix-only" (remediate, only report
|
|
# upgrades), or "notify" (never act, only report).
|
|
"autonomy": "full",
|
|
"health_interval_seconds": 300, # 5 min
|
|
"upgrade_interval_seconds": 604800, # 7 days
|
|
"matrix_poll_seconds": 30,
|
|
# A monitor must be seen DOWN this many consecutive cycles before we act,
|
|
# to ride out restarts / flapping.
|
|
"down_confirmations": 2,
|
|
# Send a quiet "all healthy" heartbeat to Matrix at most this often.
|
|
"heartbeat_interval_seconds": 86400, # daily
|
|
"state_file": "/home/paul/.local/state/semprini-maintainer/state.json",
|
|
"matrix": {
|
|
"homeserver": "https://matrix.semprini.me",
|
|
"user_id": "",
|
|
"access_token": "",
|
|
"room_id": "",
|
|
"admin_user_id": "@paul:semprini.me",
|
|
},
|
|
}
|
|
|
|
|
|
def load_config(path):
|
|
cfg = json.loads(json.dumps(DEFAULTS)) # deep copy
|
|
if path and os.path.exists(path):
|
|
with open(path) as fh:
|
|
user = json.load(fh)
|
|
for k, v in user.items():
|
|
if k == "matrix" and isinstance(v, dict):
|
|
cfg["matrix"].update(v)
|
|
else:
|
|
cfg[k] = v
|
|
return cfg
|
|
|
|
|
|
def load_state(cfg):
|
|
path = cfg["state_file"]
|
|
if os.path.exists(path):
|
|
try:
|
|
with open(path) as fh:
|
|
return json.load(fh)
|
|
except Exception:
|
|
log("WARN: could not parse state file, starting fresh")
|
|
return {
|
|
"last_health": 0,
|
|
"last_upgrade": 0,
|
|
"last_heartbeat": 0,
|
|
"down_streak": {}, # monitor name -> consecutive-down count
|
|
"matrix_since": None, # /sync pagination token
|
|
"pending": None, # outstanding question awaiting admin reply
|
|
"last_event_check": None, # max kuma heartbeat id seen (outage history)
|
|
"open_outages": {}, # monitor_id -> open DOWN window, cross-cycle
|
|
}
|
|
|
|
|
|
def save_state(cfg, state):
|
|
path = cfg["state_file"]
|
|
os.makedirs(os.path.dirname(path), exist_ok=True)
|
|
tmp = path + ".tmp"
|
|
with open(tmp, "w") as fh:
|
|
json.dump(state, fh, indent=2)
|
|
os.replace(tmp, path)
|
|
|
|
|
|
def log(msg):
|
|
print(f"[{datetime.now(timezone.utc).isoformat()}] {msg}", flush=True)
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Uptime Kuma — read tracked monitors + latest status
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
KUMA_QUERY = (
|
|
"select json_group_array(json_object("
|
|
"'id', m.id, 'name', m.name, 'type', m.type, 'url', m.url, "
|
|
"'hostname', m.hostname, 'port', m.port, "
|
|
"'status', (select h.status from heartbeat h where h.monitor_id=m.id "
|
|
"order by h.time desc limit 1), "
|
|
"'msg', (select h.msg from heartbeat h where h.monitor_id=m.id "
|
|
"order by h.time desc limit 1))) "
|
|
"from monitor m where m.active=1 order by m.id;"
|
|
)
|
|
|
|
# Uptime Kuma heartbeat status codes.
|
|
KUMA_DOWN, KUMA_UP, KUMA_PENDING, KUMA_MAINT = 0, 1, 2, 3
|
|
|
|
|
|
def read_kuma_monitors(cfg):
|
|
"""Return list of active monitors with their latest status.
|
|
|
|
Reads kuma.db out of the named Docker volume via a throwaway sqlite
|
|
container (the file is root-owned; the daemon user only has docker access).
|
|
"""
|
|
script = f"cp /src/kuma.db /tmp/k.db && sqlite3 /tmp/k.db \"{KUMA_QUERY}\""
|
|
out = subprocess.run(
|
|
[
|
|
"docker", "run", "--rm",
|
|
"-v", f"{cfg['kuma_volume']}:/src:ro",
|
|
cfg["sqlite_image"], "sh", "-c", script,
|
|
],
|
|
capture_output=True, text=True, timeout=120,
|
|
)
|
|
if out.returncode != 0:
|
|
raise RuntimeError(f"kuma read failed: {out.stderr.strip()}")
|
|
data = out.stdout.strip() or "[]"
|
|
return json.loads(data)
|
|
|
|
|
|
# Kuma flags every up/down transition with important=1. This reads those
|
|
# transitions since a given heartbeat id, plus the current max id (watermark),
|
|
# in one shot — so we can spot outages that have already self-recovered.
|
|
KUMA_EVENTS_QUERY = (
|
|
"select json_object("
|
|
"'watermark', (select max(id) from heartbeat), "
|
|
"'events', (select json_group_array(json_object("
|
|
"'monitor_id', h.monitor_id, 'name', m.name, 'status', h.status, "
|
|
"'time', h.time, 'msg', h.msg)) "
|
|
"from heartbeat h join monitor m on m.id = h.monitor_id "
|
|
"where h.important = 1 {clause} order by h.id));"
|
|
)
|
|
|
|
|
|
def read_kuma_events(cfg, since_id):
|
|
"""Return (watermark, events) of monitor up/down transitions.
|
|
|
|
`events` are the important (transition) heartbeats with id > since_id, in
|
|
chronological order; `watermark` is the current max heartbeat id to pass
|
|
back next time. since_id=None means first call: returns the watermark with
|
|
NO events (a baseline, so we don't replay weeks of history). Uses the same
|
|
throwaway sqlite container as read_kuma_monitors.
|
|
"""
|
|
clause = "and 1=0" if since_id is None else f"and h.id > {int(since_id)}"
|
|
query = KUMA_EVENTS_QUERY.format(clause=clause)
|
|
script = f"cp /src/kuma.db /tmp/k.db && sqlite3 /tmp/k.db \"{query}\""
|
|
out = subprocess.run(
|
|
[
|
|
"docker", "run", "--rm",
|
|
"-v", f"{cfg['kuma_volume']}:/src:ro",
|
|
cfg["sqlite_image"], "sh", "-c", script,
|
|
],
|
|
capture_output=True, text=True, timeout=120,
|
|
)
|
|
if out.returncode != 0:
|
|
raise RuntimeError(f"kuma events read failed: {out.stderr.strip()}")
|
|
data = json.loads(out.stdout.strip() or "{}")
|
|
return data.get("watermark"), (data.get("events") or [])
|
|
|
|
|
|
def detect_outages(events, open_outages):
|
|
"""Pair DOWN→UP transitions into self-recovered outage windows.
|
|
|
|
`open_outages` (monitor_id str -> {name, down_at, msg}) carries outages
|
|
still open from earlier cycles and is mutated in place, so an outage that
|
|
spans a cycle boundary is still paired correctly. Returns the windows that
|
|
have now recovered (a service that went DOWN and is UP again).
|
|
"""
|
|
recovered = []
|
|
for ev in events:
|
|
mid = str(ev["monitor_id"])
|
|
if ev["status"] == KUMA_DOWN:
|
|
open_outages[mid] = {"name": ev["name"], "down_at": ev["time"],
|
|
"msg": ev.get("msg")}
|
|
elif ev["status"] == KUMA_UP:
|
|
o = open_outages.pop(mid, None)
|
|
recovered.append({
|
|
"name": ev["name"],
|
|
"down_at": (o or {}).get("down_at"), # None => began before us
|
|
"up_at": ev["time"],
|
|
"msg": (o or {}).get("msg") or ev.get("msg"),
|
|
})
|
|
return recovered
|
|
|
|
|
|
def _parse_kuma_time(t):
|
|
for fmt in ("%Y-%m-%d %H:%M:%S.%f", "%Y-%m-%d %H:%M:%S"):
|
|
try:
|
|
return datetime.strptime(t, fmt)
|
|
except (ValueError, TypeError):
|
|
continue
|
|
return None
|
|
|
|
|
|
def _fmt_outage_duration(down_at, up_at):
|
|
d, u = _parse_kuma_time(down_at), _parse_kuma_time(up_at)
|
|
if not d or not u:
|
|
return "?"
|
|
secs = max(0, int((u - d).total_seconds()))
|
|
h, rem = divmod(secs, 3600)
|
|
m, s = divmod(rem, 60)
|
|
if h:
|
|
return f"{h}h{m:02d}m"
|
|
return f"{m}m{s:02d}s" if m else f"{s}s"
|
|
|
|
|
|
def docker_ps():
|
|
out = subprocess.run(
|
|
["docker", "ps", "-a", "--format",
|
|
"table {{.Names}}\t{{.Status}}\t{{.Image}}"],
|
|
capture_output=True, text=True, timeout=60,
|
|
)
|
|
return out.stdout.strip()
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Matrix client (client-server API over stdlib)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
def _matrix_req(cfg, method, path, body=None, params=None):
|
|
base = cfg["matrix"]["homeserver"].rstrip("/")
|
|
url = f"{base}{path}"
|
|
if params:
|
|
url += "?" + urllib.parse.urlencode(params)
|
|
data = json.dumps(body).encode() if body is not None else None
|
|
req = urllib.request.Request(url, data=data, method=method)
|
|
req.add_header("Authorization", f"Bearer {cfg['matrix']['access_token']}")
|
|
req.add_header("Content-Type", "application/json")
|
|
with urllib.request.urlopen(req, timeout=60) as resp:
|
|
raw = resp.read().decode()
|
|
return json.loads(raw) if raw else {}
|
|
|
|
|
|
def _md_to_html(text):
|
|
"""Very small markdown -> HTML for Matrix formatted bodies."""
|
|
import html
|
|
esc = html.escape(text)
|
|
esc = re.sub(r"\*\*(.+?)\*\*", r"<b>\1</b>", esc)
|
|
esc = re.sub(r"`(.+?)`", r"<code>\1</code>", esc)
|
|
return esc.replace("\n", "<br>")
|
|
|
|
|
|
def matrix_send(cfg, text):
|
|
if not cfg["matrix"].get("room_id") or not cfg["matrix"].get("access_token"):
|
|
log(f"MATRIX (not configured) >> {text}")
|
|
return
|
|
txn = str(int(time.time() * 1000))
|
|
room = urllib.parse.quote(cfg["matrix"]["room_id"])
|
|
try:
|
|
_matrix_req(
|
|
cfg, "PUT",
|
|
f"/_matrix/client/v3/rooms/{room}/send/m.room.message/{txn}",
|
|
body={
|
|
"msgtype": "m.text",
|
|
"body": text,
|
|
"format": "org.matrix.custom.html",
|
|
"formatted_body": _md_to_html(text),
|
|
},
|
|
)
|
|
except Exception as e:
|
|
log(f"ERROR sending to Matrix: {e}")
|
|
|
|
|
|
def matrix_poll_admin(cfg, state):
|
|
"""Return new text messages from the admin in our room since last sync."""
|
|
if not cfg["matrix"].get("access_token"):
|
|
return []
|
|
params = {"timeout": "0"}
|
|
if state.get("matrix_since"):
|
|
params["since"] = state["matrix_since"]
|
|
try:
|
|
resp = _matrix_req(cfg, "GET", "/_matrix/client/v3/sync", params=params)
|
|
except Exception as e:
|
|
log(f"ERROR polling Matrix: {e}")
|
|
return []
|
|
state["matrix_since"] = resp.get("next_batch", state.get("matrix_since"))
|
|
messages = []
|
|
room_id = cfg["matrix"]["room_id"]
|
|
admin = cfg["matrix"]["admin_user_id"]
|
|
rooms = resp.get("rooms", {}).get("join", {})
|
|
room = rooms.get(room_id)
|
|
if not room:
|
|
return messages
|
|
for ev in room.get("timeline", {}).get("events", []):
|
|
if ev.get("sender") != admin:
|
|
continue
|
|
# This bot is stdlib-only and cannot do E2E encryption. If the room is
|
|
# encrypted, the admin's replies arrive as undecryptable m.room.encrypted
|
|
# events — warn loudly rather than silently ignoring the operator. The
|
|
# room should be created unencrypted (see register-matrix-bot.sh).
|
|
if ev.get("type") == "m.room.encrypted":
|
|
log("WARN: received an ENCRYPTED message from the admin — this bot "
|
|
"cannot decrypt it. The maintenance room must be unencrypted.")
|
|
continue
|
|
if ev.get("type") != "m.room.message":
|
|
continue
|
|
content = ev.get("content", {})
|
|
if content.get("msgtype") == "m.text":
|
|
messages.append(content.get("body", "").strip())
|
|
return messages
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Stack-support agent invocation (Claude Code, headless)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
RESULT_PROTOCOL = """
|
|
When you have finished, end your reply with EXACTLY ONE fenced json block:
|
|
|
|
```json
|
|
{"status": "ok|action_taken|escalate|blocked",
|
|
"summary": "one or two sentence operator summary (markdown ok)",
|
|
"details": "optional longer detail",
|
|
"question": "optional — only if status is escalate: the decision you need",
|
|
"options": ["optional", "list", "of", "choices"]}
|
|
```
|
|
|
|
status meanings:
|
|
ok — everything healthy, nothing to do
|
|
action_taken — you fixed/changed something; describe it in summary
|
|
escalate — you need a human decision (set question + options)
|
|
blocked — hard blocker, cannot proceed safely (explain in summary)
|
|
""".strip()
|
|
|
|
|
|
def run_agent(cfg, task):
|
|
"""Invoke the canonical stack-support agent headless and return its output."""
|
|
agent_path = os.path.join(cfg["repo_dir"], cfg["agent_file"])
|
|
with open(agent_path) as fh:
|
|
agent_md = fh.read()
|
|
|
|
system = (
|
|
"You are the semprini-core Stack Maintenance Agent, running unattended "
|
|
"as a systemd service. The following is your canonical agent definition; "
|
|
"follow it as the source of truth, together with "
|
|
f"{cfg['runbook_file']} and CLAUDE.md in this repository.\n\n"
|
|
"You operate FULLY AUTONOMOUSLY per the CLAUDE.md sysadmin role: own the "
|
|
"fix end to end. Make the smallest safe change that restores service. "
|
|
"Pin versions after successful upgrades; auto-rollback on failure. "
|
|
"Only set status=escalate when a genuine human decision is required "
|
|
"(e.g. a major/breaking upgrade with trade-offs).\n\n"
|
|
"Services in scope are exactly those Uptime Kuma tracks. Other projects "
|
|
"(e.g. semprini-blog) live alongside semprini-core and use it for "
|
|
"IAM/reverse-proxy; maintain their containers too when they are tracked.\n\n"
|
|
"--- BEGIN agents/stack-support.agent.md ---\n"
|
|
f"{agent_md}\n"
|
|
"--- END agents/stack-support.agent.md ---\n\n"
|
|
f"{RESULT_PROTOCOL}"
|
|
)
|
|
|
|
cmd = [
|
|
cfg["claude_bin"], "-p", task,
|
|
"--append-system-prompt", system,
|
|
"--permission-mode", "bypassPermissions",
|
|
"--output-format", "text",
|
|
]
|
|
if cfg.get("claude_model"):
|
|
cmd += ["--model", cfg["claude_model"]]
|
|
for d in cfg.get("extra_dirs", []):
|
|
cmd += ["--add-dir", d]
|
|
cmd += list(cfg.get("claude_extra_args", []))
|
|
|
|
log(f"Invoking stack-support agent ({len(task)} char task)…")
|
|
proc = subprocess.run(
|
|
cmd, cwd=cfg["repo_dir"], capture_output=True, text=True,
|
|
timeout=cfg["agent_timeout_seconds"],
|
|
)
|
|
if proc.returncode != 0:
|
|
log(f"WARN: agent exited {proc.returncode}: {proc.stderr.strip()[:500]}")
|
|
return proc.stdout.strip()
|
|
|
|
|
|
def parse_agent_result(output):
|
|
"""Pull the trailing JSON envelope out of the agent's reply."""
|
|
blocks = re.findall(r"```json\s*(\{.*?\})\s*```", output, re.DOTALL)
|
|
if blocks:
|
|
try:
|
|
return json.loads(blocks[-1])
|
|
except Exception:
|
|
pass
|
|
# Fallback: try a bare trailing {...}
|
|
m = re.search(r"(\{[^{}]*\"status\"[^{}]*\})\s*$", output, re.DOTALL)
|
|
if m:
|
|
try:
|
|
return json.loads(m.group(1))
|
|
except Exception:
|
|
pass
|
|
return {"status": "unknown",
|
|
"summary": (output[-800:] or "no output from agent")}
|
|
|
|
|
|
def relay_result(cfg, state, header, result):
|
|
"""Post an agent result to Matrix and register any escalation."""
|
|
status = result.get("status", "unknown")
|
|
summary = result.get("summary", "")
|
|
icon = {"ok": "✅", "action_taken": "🔧", "escalate": "❓",
|
|
"blocked": "⛔", "unknown": "⚠️"}.get(status, "•")
|
|
msg = f"{icon} **{header}** — {status}\n{summary}"
|
|
if result.get("details"):
|
|
msg += f"\n\n{result['details']}"
|
|
|
|
if status == "escalate" and result.get("question"):
|
|
opts = result.get("options") or []
|
|
msg += f"\n\n**Decision needed:** {result['question']}"
|
|
if opts:
|
|
msg += "\n" + "\n".join(f" {i+1}. {o}" for i, o in enumerate(opts))
|
|
msg += "\n\n_Reply with a number or a short answer._"
|
|
state["pending"] = {
|
|
"header": header,
|
|
"question": result["question"],
|
|
"options": opts,
|
|
"asked_at": time.time(),
|
|
}
|
|
matrix_send(cfg, msg)
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Cycles
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
def status_label(code):
|
|
return {KUMA_DOWN: "DOWN", KUMA_UP: "UP", KUMA_PENDING: "PENDING",
|
|
KUMA_MAINT: "MAINTENANCE"}.get(code, f"?{code}")
|
|
|
|
|
|
def report_recovered_outages(cfg, state, recovered):
|
|
"""A tracked service went DOWN then recovered on its own. Surface the
|
|
window and, unless notify-only, have the agent dig into logs for that
|
|
period so recurring faults get caught even though status is now UP."""
|
|
def fmt(o):
|
|
when = o["down_at"] or "(began before last check)"
|
|
dur = _fmt_outage_duration(o["down_at"], o["up_at"]) if o["down_at"] else "?"
|
|
return (f" • **{o['name']}**: DOWN {when} → recovered {o['up_at']} "
|
|
f"({dur}) — {o.get('msg') or 'n/a'}")
|
|
lines = "\n".join(fmt(o) for o in recovered)
|
|
|
|
if cfg["autonomy"] == "notify":
|
|
matrix_send(cfg, "🟡 **Self-recovered outage(s)** (notify-only — not "
|
|
f"investigated):\n{lines}")
|
|
return
|
|
|
|
task = (
|
|
"POST-INCIDENT REVIEW — one or more tracked services went DOWN and then "
|
|
"recovered on their own since the last check. Current status is UP, so "
|
|
"no remediation is needed; investigate the likely cause so recurring "
|
|
"faults are caught early.\n\n"
|
|
f"Self-recovered outage windows:\n{lines}\n\n"
|
|
"Current containers:\n"
|
|
f"{docker_ps()}\n\n"
|
|
"For each: map the monitor target to its container, inspect logs and "
|
|
"events around the outage window (e.g. `docker logs --since <down> "
|
|
"--until <up>`, container restart counts, journald), and identify the "
|
|
"probable cause (OOM, restart/redeploy, dependency flap, network). Do "
|
|
"NOT change anything unless something is actively still wrong. Report "
|
|
"concise per-service findings; set status=escalate only if a fix needs "
|
|
"a human decision."
|
|
)
|
|
result = parse_agent_result(run_agent(cfg, task))
|
|
relay_result(cfg, state, "Outage post-mortem", result)
|
|
|
|
|
|
def health_cycle(cfg, state):
|
|
log("Health cycle: reading Uptime Kuma…")
|
|
monitors = read_kuma_monitors(cfg)
|
|
streak = state.setdefault("down_streak", {})
|
|
|
|
down_now = []
|
|
for m in monitors:
|
|
name = m["name"]
|
|
if m.get("status") == KUMA_DOWN:
|
|
streak[name] = streak.get(name, 0) + 1
|
|
down_now.append(m)
|
|
else:
|
|
streak.pop(name, None)
|
|
|
|
confirmed = [m for m in down_now
|
|
if streak.get(m["name"], 0) >= cfg["down_confirmations"]]
|
|
|
|
log(f" {len(monitors)} tracked, {len(down_now)} down, "
|
|
f"{len(confirmed)} confirmed down")
|
|
|
|
# Surface services that blipped DOWN→UP since we last looked. Current status
|
|
# may be fine, but a self-recovered outage is a pointer worth investigating.
|
|
try:
|
|
open_outages = state.setdefault("open_outages", {})
|
|
watermark, events = read_kuma_events(cfg, state.get("last_event_check"))
|
|
recovered = detect_outages(events, open_outages)
|
|
state["last_event_check"] = watermark
|
|
if recovered:
|
|
log(f" {len(recovered)} self-recovered outage(s) since last check")
|
|
report_recovered_outages(cfg, state, recovered)
|
|
except Exception as e:
|
|
log(f"WARN: outage-history check failed: {e}")
|
|
|
|
if not confirmed:
|
|
# quiet heartbeat
|
|
if (time.time() - state.get("last_heartbeat", 0)
|
|
>= cfg["heartbeat_interval_seconds"]):
|
|
up = sum(1 for m in monitors if m.get("status") == KUMA_UP)
|
|
matrix_send(
|
|
cfg,
|
|
f"✅ **Daily health** — {up}/{len(monitors)} tracked services up.",
|
|
)
|
|
state["last_heartbeat"] = time.time()
|
|
return
|
|
|
|
if cfg["autonomy"] == "notify":
|
|
lines = "\n".join(
|
|
f" • **{m['name']}** ({m.get('url') or m.get('hostname')}): "
|
|
f"{m.get('msg') or 'down'}" for m in confirmed
|
|
)
|
|
matrix_send(cfg, f"⚠️ **Services down** (notify-only mode):\n{lines}")
|
|
return
|
|
|
|
incident = "\n".join(
|
|
f" • {m['name']} | type={m['type']} | "
|
|
f"target={m.get('url') or (m.get('hostname') or '') + ':' + str(m.get('port') or '')} | "
|
|
f"last_msg={m.get('msg') or 'n/a'}"
|
|
for m in confirmed
|
|
)
|
|
task = (
|
|
"ROUTINE HEALTH CHECK — Uptime Kuma reports the following tracked "
|
|
"services DOWN (confirmed over multiple cycles):\n\n"
|
|
f"{incident}\n\n"
|
|
"Current containers:\n"
|
|
f"{docker_ps()}\n\n"
|
|
"For each down service: map the monitor target (a compose service "
|
|
"name / network alias) to its container, diagnose using the runbook, "
|
|
"and remediate autonomously. Re-verify recovery. If a service is down "
|
|
"by design or needs a human decision, set status=escalate."
|
|
)
|
|
result = parse_agent_result(run_agent(cfg, task))
|
|
relay_result(cfg, state, "Health remediation", result)
|
|
|
|
|
|
def resolve_tracked_images(cfg, monitors):
|
|
"""Best-effort map of tracked services -> running image (for context)."""
|
|
ps = subprocess.run(
|
|
["docker", "ps", "--format", "{{.Names}}\t{{.Image}}"],
|
|
capture_output=True, text=True, timeout=60,
|
|
).stdout.strip().splitlines()
|
|
name_to_image = {}
|
|
for line in ps:
|
|
if "\t" in line:
|
|
n, img = line.split("\t", 1)
|
|
name_to_image[n] = img
|
|
rows = []
|
|
for m in monitors:
|
|
target = (m.get("hostname") or m.get("url") or "").lower()
|
|
token = re.sub(r"^https?://", "", target).split("/")[0].split(":")[0]
|
|
matches = [f"{n} -> {img}" for n, img in name_to_image.items()
|
|
if token and token in n]
|
|
rows.append(f" • {m['name']} (target={token or 'n/a'}): "
|
|
f"{', '.join(matches) or 'no obvious container match'}")
|
|
return "\n".join(rows)
|
|
|
|
|
|
def upgrade_cycle(cfg, state):
|
|
if cfg["autonomy"] == "notify":
|
|
autonomy_note = (
|
|
"Mode is NOTIFY-ONLY: do NOT change anything. Report available "
|
|
"upgrades and recommendations only; set status=escalate with the list."
|
|
)
|
|
elif cfg["autonomy"] == "fix-only":
|
|
autonomy_note = (
|
|
"Mode is FIX-ONLY for upgrades: research and report available "
|
|
"upgrades; do NOT apply them. Set status=escalate to ask which to apply."
|
|
)
|
|
else:
|
|
autonomy_note = (
|
|
"Mode is FULL: apply safe upgrades autonomously (pull, test for "
|
|
"stability, pin the new tag in the owning compose file, redeploy, "
|
|
"rollback on failure). For MAJOR/breaking upgrades with trade-offs, "
|
|
"set status=escalate with options instead of applying."
|
|
)
|
|
|
|
log("Upgrade cycle…")
|
|
monitors = read_kuma_monitors(cfg)
|
|
context = resolve_tracked_images(cfg, monitors)
|
|
task = (
|
|
"PERIODIC UPGRADE REVIEW. Scope = services Uptime Kuma currently tracks. "
|
|
"Tracked services and their running images:\n\n"
|
|
f"{context}\n\n"
|
|
"Use the upgrade-services skill (agents/skills/upgrade-services/SKILL.md). "
|
|
"For each tracked service: find the compose file that owns it (search "
|
|
"core_stack/compose*.yml and other project repos under the added dirs, "
|
|
"e.g. semprini-blog/docker-compose.yml), determine the current pinned "
|
|
"tag, check Docker Hub / GitHub for newer stable releases, and assess "
|
|
"breaking changes.\n\n"
|
|
f"{autonomy_note}\n\n"
|
|
"Summarise what is up to date, what you upgraded, and anything needing "
|
|
"a decision."
|
|
)
|
|
result = parse_agent_result(run_agent(cfg, task))
|
|
relay_result(cfg, state, "Upgrade review", result)
|
|
|
|
|
|
def handle_replies(cfg, state):
|
|
msgs = matrix_poll_admin(cfg, state)
|
|
if not msgs:
|
|
return
|
|
pending = state.get("pending")
|
|
for body in msgs:
|
|
log(f"Admin reply: {body!r}")
|
|
if not pending:
|
|
# Unprompted instruction from the admin — treat as an ad-hoc task.
|
|
task = (
|
|
"The operator sent this instruction over Matrix; act on it "
|
|
f"autonomously and report back:\n\n{body}"
|
|
)
|
|
result = parse_agent_result(run_agent(cfg, task))
|
|
relay_result(cfg, state, "Operator request", result)
|
|
pending = state.get("pending")
|
|
continue
|
|
|
|
# Resolve the outstanding escalation with the admin's answer.
|
|
answer = body
|
|
opts = pending.get("options") or []
|
|
if body.isdigit() and 1 <= int(body) <= len(opts):
|
|
answer = opts[int(body) - 1]
|
|
task = (
|
|
f"Earlier you escalated this decision:\n"
|
|
f" Question: {pending['question']}\n"
|
|
f" Options: {opts}\n\n"
|
|
f"The operator decided: {answer}\n\n"
|
|
"Carry out that decision now, autonomously, with full testing and "
|
|
"rollback on failure. Report the outcome."
|
|
)
|
|
state["pending"] = None
|
|
save_state(cfg, state)
|
|
result = parse_agent_result(run_agent(cfg, task))
|
|
relay_result(cfg, state, "Decision carried out", result)
|
|
pending = state.get("pending")
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Main loop
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
def main():
|
|
ap = argparse.ArgumentParser(description="semprini stack maintenance daemon")
|
|
ap.add_argument("--config",
|
|
default="/home/paul/Dev/semprini-maintainer/config.json")
|
|
ap.add_argument("--once", choices=["health", "upgrade", "replies"],
|
|
help="run a single cycle and exit (for testing)")
|
|
args = ap.parse_args()
|
|
|
|
cfg = load_config(args.config)
|
|
state = load_state(cfg)
|
|
|
|
if args.once:
|
|
log(f"Running single '{args.once}' cycle…")
|
|
{"health": health_cycle, "upgrade": upgrade_cycle,
|
|
"replies": handle_replies}[args.once](cfg, state)
|
|
save_state(cfg, state)
|
|
return
|
|
|
|
log("semprini-maintainer starting up.")
|
|
matrix_send(cfg, "🤖 **stack-maintainer online** — monitoring tracked "
|
|
"services via Uptime Kuma.")
|
|
|
|
while True:
|
|
now = time.time()
|
|
try:
|
|
handle_replies(cfg, state)
|
|
|
|
if now - state.get("last_health", 0) >= cfg["health_interval_seconds"]:
|
|
health_cycle(cfg, state)
|
|
state["last_health"] = now
|
|
|
|
# Hold upgrades while an escalation is unanswered.
|
|
if (not state.get("pending")
|
|
and now - state.get("last_upgrade", 0)
|
|
>= cfg["upgrade_interval_seconds"]):
|
|
upgrade_cycle(cfg, state)
|
|
state["last_upgrade"] = now
|
|
|
|
save_state(cfg, state)
|
|
except Exception as e:
|
|
log("ERROR in main loop:\n" + traceback.format_exc())
|
|
try:
|
|
matrix_send(cfg, f"⛔ **maintainer error**: {e}")
|
|
except Exception:
|
|
pass
|
|
|
|
time.sleep(cfg["matrix_poll_seconds"])
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|