#!/usr/bin/env python3 """Stateful external production health and operations notifications.""" from __future__ import annotations import argparse import json import os import re import smtplib import ssl import sys import tempfile import urllib.error import urllib.request from datetime import datetime, timezone from email.message import EmailMessage from pathlib import Path from typing import Any MONITORED_COUNTERS = { "who_need_help_http_exceptions_total", "who_need_help_oban_jobs_failed_total", "who_need_help_email_deliveries_total", "who_need_help_email_delivery_exceptions_total", } PROMETHEUS_SAMPLE = re.compile( r"^(?P[a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{(?P.*)\})?\s+" r"(?P[-+]?(?:[0-9]+(?:\.[0-9]*)?|\.[0-9]+)(?:[eE][-+]?[0-9]+)?)" r"(?:\s+[0-9]+)?$" ) PROMETHEUS_LABEL = re.compile(r'([a-zA-Z_][a-zA-Z0-9_]*)="((?:\\.|[^"\\])*)"') def utc_now() -> str: return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") def read_json(path: Path) -> dict[str, Any]: with path.open("r", encoding="utf-8") as handle: value = json.load(handle) if not isinstance(value, dict): raise ValueError(f"Expected a JSON object in {path}") return value def write_json_atomic(path: Path, value: dict[str, Any]) -> None: path.parent.mkdir(mode=0o700, parents=True, exist_ok=True) descriptor, temporary_name = tempfile.mkstemp(dir=path.parent, prefix=f".{path.name}.") temporary = Path(temporary_name) try: with os.fdopen(descriptor, "w", encoding="utf-8") as handle: json.dump(value, handle, ensure_ascii=False, indent=2, sort_keys=True) handle.write("\n") temporary.chmod(0o600) temporary.replace(path) finally: temporary.unlink(missing_ok=True) def required_text(mapping: dict[str, Any], name: str) -> str: value = mapping.get(name) if not isinstance(value, str) or not value.strip(): raise ValueError(f"Missing non-empty configuration value: {name}") return value.strip() def parse_labels(raw_labels: str | None) -> dict[str, str]: if not raw_labels: return {} labels: dict[str, str] = {} position = 0 for match in PROMETHEUS_LABEL.finditer(raw_labels): separator = raw_labels[position : match.start()].strip() if separator not in {"", ","}: raise ValueError("Malformed Prometheus labels") labels[match.group(1)] = bytes(match.group(2), "utf-8").decode("unicode_escape") position = match.end() if raw_labels[position:].strip() not in {"", ","}: raise ValueError("Malformed Prometheus labels") return labels def parse_monitored_counters(payload: str) -> dict[str, float]: counters: dict[str, float] = {} for raw_line in payload.splitlines(): line = raw_line.strip() if not line or line.startswith("#"): continue match = PROMETHEUS_SAMPLE.fullmatch(line) if not match or match.group("name") not in MONITORED_COUNTERS: continue name = match.group("name") labels = parse_labels(match.group("labels")) if name == "who_need_help_email_deliveries_total": if labels.get("status") != "error": continue key = f"{name}|status=error" elif name == "who_need_help_oban_jobs_failed_total": key = f"{name}|queue={labels.get('queue', 'unknown')}" else: key = name counters[key] = counters.get(key, 0.0) + float(match.group("value")) return counters def counter_increases( previous: dict[str, Any] | None, current: dict[str, float] ) -> dict[str, float]: if previous is None: return {} increases: dict[str, float] = {} for key, current_value in current.items(): previous_value = float(previous.get(key, 0.0)) if current_value > previous_value: increases[key] = current_value - previous_value return increases def send_message(config: dict[str, Any], subject: str, body: str) -> None: smtp = config.get("smtp") if not isinstance(smtp, dict): raise ValueError("Missing SMTP configuration") relay = required_text(smtp, "relay") port = int(smtp.get("port")) username = required_text(smtp, "username") password = required_text(smtp, "password") from_address = required_text(smtp, "from_address") from_name = required_text(smtp, "from_name") recipient = required_text(smtp, "recipient") implicit_ssl = bool(smtp.get("implicit_ssl", False)) starttls = bool(smtp.get("starttls", True)) message = EmailMessage() message["From"] = f"{from_name} <{from_address}>" message["To"] = recipient message["Subject"] = subject message.set_content(body) context = ssl.create_default_context() if implicit_ssl: client_context = smtplib.SMTP_SSL(relay, port, timeout=30, context=context) else: client_context = smtplib.SMTP(relay, port, timeout=30) with client_context as client: if not implicit_ssl: client.ehlo() if starttls: client.starttls(context=context) client.ehlo() client.login(username, password) client.send_message(message) def check_health(config: dict[str, Any]) -> tuple[str, str]: url = required_text(config, "health_url") timeout = float(config.get("health_timeout_seconds")) request = urllib.request.Request( url, headers={"User-Agent": "WhoNeedHelp-External-Monitor/1.0"}, ) try: with urllib.request.urlopen(request, timeout=timeout) as response: status = response.status payload = response.read(4096) parsed = json.loads(payload.decode("utf-8")) if status == 200 and parsed == {"status": "ready"}: return "up", f"HTTP {status}; ready payload matched" return "down", f"HTTP {status}; unexpected readiness payload" except (OSError, ValueError, urllib.error.URLError) as error: return "down", f"{type(error).__name__}: {str(error)[:500]}" def check_metrics(config: dict[str, Any]) -> tuple[str, str, str, dict[str, float]]: url = required_text(config, "metrics_url") token = required_text(config, "metrics_token") timeout = float(config.get("metrics_timeout_seconds")) request = urllib.request.Request( url, headers={ "Authorization": f"Bearer {token}", "User-Agent": "WhoNeedHelp-External-Monitor/1.0", }, ) try: with urllib.request.urlopen(request, timeout=timeout) as response: status = response.status node = response.headers.get("x-wnh-node", "unknown")[:200] payload = response.read(4 * 1024 * 1024 + 1) if status != 200: return "down", f"HTTP {status}", node, {} if len(payload) > 4 * 1024 * 1024: return "down", "metrics response exceeded 4 MiB", node, {} counters = parse_monitored_counters(payload.decode("utf-8")) return "up", f"HTTP {status}; aggregate counters parsed", node, counters except (OSError, UnicodeError, ValueError, urllib.error.URLError) as error: return "down", f"{type(error).__name__}: {str(error)[:500]}", "unknown", {} def monitor(config_path: Path, state_path: Path) -> int: config = read_json(config_path) previous: dict[str, Any] = {} if state_path.exists(): previous = read_json(state_path) health_status, health_detail = check_health(config) metrics_status, metrics_detail, metrics_node, counters = check_metrics(config) status = "up" if health_status == "up" and metrics_status == "up" else "down" detail = f"readiness={health_status} ({health_detail}); metrics={metrics_status} ({metrics_detail})" previous_status = previous.get("status") if status != previous_status: if status == "down": send_message( config, "[Who Need Help] Production monitor is DOWN", "\n".join( [ "The independent production readiness or metrics check failed.", f"URL: {required_text(config, 'health_url')}", f"Metrics URL: {required_text(config, 'metrics_url')}", f"Observed at: {utc_now()}", f"Result: {detail}", "This alert is sent once per state transition.", ] ), ) elif previous_status == "down": send_message( config, "[Who Need Help] Production monitor recovered", "\n".join( [ "The independent production readiness and metrics checks recovered.", f"URL: {required_text(config, 'health_url')}", f"Metrics URL: {required_text(config, 'metrics_url')}", f"Observed at: {utc_now()}", f"Result: {detail}", ] ), ) previous_by_node = previous.get("metrics_by_node", {}) if not isinstance(previous_by_node, dict): previous_by_node = {} previous_counters = previous_by_node.get(metrics_node) if not isinstance(previous_counters, dict): previous_counters = None increases = counter_increases(previous_counters, counters) if metrics_status == "up" else {} if increases: lines = [ "One or more aggregate production failure counters increased.", f"Metrics node: {metrics_node}", f"Observed at: {utc_now()}", "Newly observed failures:", ] lines.extend(f"- {key}: +{value:g}" for key, value in sorted(increases.items())) lines.append("No recipient address, message content, or request payload is included.") send_message( config, "[Who Need Help] Production failure counters increased", "\n".join(lines), ) metrics_by_node = dict(previous_by_node) if metrics_status == "up": metrics_by_node[metrics_node] = counters write_json_atomic( state_path, { "checked_at": utc_now(), "detail": detail, "health_status": health_status, "metrics_by_node": metrics_by_node, "metrics_node": metrics_node, "metrics_status": metrics_status, "status": status, }, ) print(json.dumps({"status": status, "detail": detail}, sort_keys=True)) return 0 if status == "up" else 1 def notify_backup_failure(config_path: Path, unit: str) -> int: config = read_json(config_path) send_message( config, "[Who Need Help] Production backup or restore verification failed", "\n".join( [ "The scheduled encrypted production backup did not finish successfully.", f"Unit: {unit}", f"Observed at: {utc_now()}", "Inspect the local user-systemd journal and do not treat the newest snapshot as verified until a restore drill passes.", ] ), ) print("Backup failure notification sent.") return 0 def send_test_notification(config_path: Path) -> int: config = read_json(config_path) send_message( config, "[Who Need Help] Operations monitoring test", "\n".join( [ "This is a one-time delivery verification for the independent production monitor.", f"Health URL: {required_text(config, 'health_url')}", f"Metrics URL: {required_text(config, 'metrics_url')}", f"Sent at: {utc_now()}", "No production incident was detected and no application data was changed.", ] ), ) print("Operations monitoring test notification sent.") return 0 def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser() parser.add_argument( "--config", type=Path, default=Path.home() / ".config/who-need-help/monitor.json", ) parser.add_argument( "--state", type=Path, default=Path.home() / ".local/state/who-need-help/monitor.json", ) subparsers = parser.add_subparsers(dest="action", required=True) subparsers.add_parser("check") subparsers.add_parser("send-test-notification") backup = subparsers.add_parser("notify-backup-failure") backup.add_argument("--unit", required=True) return parser.parse_args() def main() -> int: args = parse_args() try: if args.action == "check": return monitor(args.config, args.state) if args.action == "send-test-notification": return send_test_notification(args.config) return notify_backup_failure(args.config, args.unit) except (OSError, ValueError, smtplib.SMTPException, json.JSONDecodeError) as error: print(f"Operations monitor failed: {type(error).__name__}: {error}", file=sys.stderr) return 1 if __name__ == "__main__": raise SystemExit(main())