373 lines
13 KiB
Python
Executable File
373 lines
13 KiB
Python
Executable File
#!/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<name>[a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{(?P<labels>.*)\})?\s+"
|
|
r"(?P<value>[-+]?(?:[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())
|