#!/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", "who_need_help_email_by_kind_deliveries_total", } MONITORED_EMAIL_KINDS = { "auth_email_change", "auth_google_link", "auth_login", "auth_registration", "content_removal_confirmation", "content_removal_operator", "content_removal_received", "content_removal_update", "support_confirmation", "support_operator", "support_update", } 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_email_by_kind_deliveries_total": kind = labels.get("kind") if labels.get("status") not in {"error", "exception"}: continue if kind not in MONITORED_EMAIL_KINDS: continue key = f"{name}|kind={kind}|status={labels['status']}" 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 parse_utc_timestamp(value: Any) -> datetime: if not isinstance(value, str) or not value.strip(): raise ValueError("Missing verified backup timestamp") parsed = datetime.fromisoformat(value.strip().replace("Z", "+00:00")) if parsed.tzinfo is None: raise ValueError("Backup timestamp must include a UTC offset") return parsed.astimezone(timezone.utc) def check_backup_freshness( config: dict[str, Any], *, now: datetime | None = None ) -> tuple[str, str]: backup = config.get("backup") if backup is None: return "disabled", "backup freshness monitoring is not configured" if not isinstance(backup, dict): raise ValueError("Backup monitoring configuration must be an object") heartbeat_path = Path(required_text(backup, "heartbeat_path")) if not heartbeat_path.is_absolute(): raise ValueError("Backup heartbeat path must be absolute") raw_max_age = backup.get("max_age_seconds") if isinstance(raw_max_age, bool): raise ValueError("Backup max_age_seconds must be a positive integer") try: max_age_seconds = int(raw_max_age) except (TypeError, ValueError) as error: raise ValueError("Backup max_age_seconds must be a positive integer") from error if max_age_seconds <= 0: raise ValueError("Backup max_age_seconds must be positive") try: heartbeat = read_json(heartbeat_path) if heartbeat.get("restore_verified") is not True: return "down", "latest backup heartbeat is not restore-verified" required_text(heartbeat, "snapshot_id") verified_at = parse_utc_timestamp(heartbeat.get("verified_at")) except (OSError, ValueError, json.JSONDecodeError) as error: return "down", f"{type(error).__name__}: {str(error)[:500]}" observed_at = (now or datetime.now(timezone.utc)).astimezone(timezone.utc) if verified_at > observed_at: return "down", "latest backup heartbeat timestamp is in the future" age_seconds = max(0, int((observed_at - verified_at).total_seconds())) if age_seconds > max_age_seconds: return ( "down", f"latest restore-verified backup is stale (age={age_seconds}s; max={max_age_seconds}s)", ) return ( "up", f"latest restore-verified backup is fresh (age={age_seconds}s; max={max_age_seconds}s)", ) 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) backup_status, backup_detail = check_backup_freshness(config) status = ( "up" if health_status == "up" and metrics_status == "up" and backup_status in {"up", "disabled"} else "down" ) detail = "; ".join( [ f"readiness={health_status} ({health_detail})", f"metrics={metrics_status} ({metrics_detail})", f"backup={backup_status} ({backup_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, metrics, or backup-freshness 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, metrics, and configured backup-freshness 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(), "backup_status": backup_status, "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())