From e779188e8e92d52665688a8367ddfb8b20081b2b Mon Sep 17 00:00:00 2001 From: SimpleTest Date: Sun, 9 Aug 2026 07:03:45 +0300 Subject: [PATCH] Monitor aggregate production delivery failures --- docs/operations.md | 52 ++++-- lib/who_need_help_web/telemetry.ex | 14 ++ .../install-production-external-monitor.sh | 27 ++- scripts/production-external-monitor.py | 151 +++++++++++++++- scripts/quality.sh | 4 + .../production_external_monitor_test.py | 170 ++++++++++++++++++ .../controllers/metrics_controller_test.exs | 44 +++++ 7 files changed, 433 insertions(+), 29 deletions(-) create mode 100644 test/scripts/production_external_monitor_test.py diff --git a/docs/operations.md b/docs/operations.md index d9c97c6..6eaf70f 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -700,13 +700,27 @@ copy it into the backup repository. There is intentionally no automatic `forget`/`prune` policy yet: retention, RPO, RTO, capacity, and key custody are operator decisions, and the script does not invent them. -## Independent public readiness monitor +## Independent production monitor The external monitor runs on an SSH target that must resolve to a host other -than production. It requests only the public readiness endpoint and accepts -only HTTP 200 with the exact JSON object `{"status":"ready"}`. State is kept -in a mode-`0600` JSON file. Email is sent once when the state changes to down -and once when it recovers; repeated down checks do not send repeated alerts. +than production. It requests the public readiness endpoint and accepts only +HTTP 200 with the exact JSON object `{"status":"ready"}`. It also scrapes the +authenticated `/metrics` endpoint with the independent `METRICS_TOKEN` copied +from the production environment into the external host's mode-`0600` +configuration. Readiness or metrics-scrape failure changes the monitor state to +down. Email is sent once when that state changes to down and once when it +recovers; repeated down checks do not send repeated alerts. + +For each observed application node, the monitor keeps a baseline of aggregate +HTTP exception, failed Oban job, and failed/exceptional email-delivery counters. +It sends one aggregate notification when one of those counters increases. A +new node establishes a baseline without alerting, and a lower value is treated +as a counter reset rather than a failure. The metric labels are deliberately +bounded to queue and delivery status; recipient addresses, message content, +request payloads, and user identifiers are never included. With multiple +application replicas, this external check samples only the node that answers +each request; use the full Prometheus deployment when every replica must be +scraped continuously. Install it from the production SMTP configuration without printing the SMTP credential: @@ -719,13 +733,14 @@ ssh buyvm-maya \ 'journalctl --user -u who-need-help-production-monitor.service --no-pager' ``` -The default schedule is once per minute with a three-second HTTP timeout. Those -defaults match the current application container health timeout and were -installed only after public readiness requests from the selected monitor host -were observed to finish below one second. They are observations of the current -path, not universal capacity or availability guarantees. Re-running the -installer refreshes the exact script, mode-`0600` configuration, and user -units, then performs one check before enabling the timer. +The default schedule is once per minute with separate three-second readiness +and metrics timeouts. Those defaults match the current application container +health timeout and were installed only after public readiness requests from +the selected monitor host were observed to finish below one second. They are +observations of the current path, not universal capacity or availability +guarantees. Re-running the installer refreshes the exact script, mode-`0600` +configuration, and user units, then performs one check before enabling the +timer. ## Local external-service boundary drill @@ -1028,11 +1043,14 @@ curl --fail \ The endpoint returns `401` without the exact token, disables response caching, and does not put the credential in a URL. The reporter exports cumulative HTTP request and duration, router exception, database query and duration, WebSocket -connection, VM memory, and scheduler run-queue metrics. Cumulative durations are -integer microseconds because the selected reporter's sum accumulator is -integer-based; divide by `1_000_000` in PromQL when seconds are required. -Definitions intentionally have no request path, user, request, or event-name -labels that could create unbounded cardinality. +connection, Oban job, aggregate single-email delivery outcome, VM memory, and +scheduler run-queue metrics. Email metrics distinguish only `ok` and `error` +adapter results and count raised delivery exceptions separately. Cumulative +durations are integer microseconds because the selected reporter's sum +accumulator is integer-based; divide by `1_000_000` in PromQL when seconds are +required. Definitions intentionally have no request path, user, request, +recipient, message, or event-name labels that could create unbounded +cardinality or expose private data. The locked `telemetry_metrics_prometheus_core` reporter aggregates distribution samples only when a scrape occurs. Each application VM therefore diff --git a/lib/who_need_help_web/telemetry.ex b/lib/who_need_help_web/telemetry.ex index 14eb68f..529b549 100644 --- a/lib/who_need_help_web/telemetry.ex +++ b/lib/who_need_help_web/telemetry.ex @@ -157,6 +157,20 @@ defmodule WhoNeedHelpWeb.Telemetry do tags: [:queue], description: "Failed Oban job attempts" ), + counter("who_need_help.email.deliveries.total", + event_name: [:swoosh, :deliver, :stop], + measurement: :duration, + tags: [:status], + tag_values: fn metadata -> + %{status: if(Map.get(metadata, :error), do: "error", else: "ok")} + end, + description: "Completed single-email delivery attempts by adapter result" + ), + counter("who_need_help.email.delivery.exceptions.total", + event_name: [:swoosh, :deliver, :exception], + measurement: :duration, + description: "Single-email delivery attempts that raised an exception" + ), sum("who_need_help.oban.job.duration.microseconds.total", event_name: [:oban, :job, :stop], measurement: fn measurements -> diff --git a/scripts/install-production-external-monitor.sh b/scripts/install-production-external-monitor.sh index b243eb3..152bdb6 100755 --- a/scripts/install-production-external-monitor.sh +++ b/scripts/install-production-external-monitor.sh @@ -10,8 +10,10 @@ remote_root=${MONITOR_REMOTE_ROOT:-/home/simple/.local/lib/who-need-help} remote_config=${MONITOR_REMOTE_CONFIG:-/home/simple/.config/who-need-help/monitor.json} remote_state=${MONITOR_REMOTE_STATE:-/home/simple/.local/state/who-need-help/monitor.json} monitor_url=${MONITOR_URL:-https://whoneedhelp.com/healthz/ready} +metrics_url=${MONITOR_METRICS_URL:-https://whoneedhelp.com/metrics} monitor_calendar=${MONITOR_ON_CALENDAR:-'*:0/1'} health_timeout=${MONITOR_HEALTH_TIMEOUT_SECONDS:-3} +metrics_timeout=${MONITOR_METRICS_TIMEOUT_SECONDS:-3} for command in mktemp scp ssh; do command -v "$command" >/dev/null 2>&1 || { @@ -32,6 +34,11 @@ case "$health_timeout" in 0) echo "MONITOR_HEALTH_TIMEOUT_SECONDS must be a positive integer." >&2; exit 2 ;; esac +case "$metrics_timeout" in + '' | *[!0-9]*) echo "MONITOR_METRICS_TIMEOUT_SECONDS must be a positive integer." >&2; exit 2 ;; + 0) echo "MONITOR_METRICS_TIMEOUT_SECONDS must be a positive integer." >&2; exit 2 ;; +esac + work_dir=$(mktemp -d "$ROOT/tmp/production-operations/monitor-install.XXXXXX") config="$work_dir/monitor.json" service="$work_dir/who-need-help-production-monitor.service" @@ -43,12 +50,12 @@ cleanup() { trap cleanup EXIT HUP INT TERM ssh -o BatchMode=yes "$source_target" \ - "python3 - '$production_env' '$monitor_url' '$health_timeout'" >"$config" <<'PY' + "python3 - '$production_env' '$monitor_url' '$metrics_url' '$health_timeout' '$metrics_timeout'" >"$config" <<'PY' import json import shlex import sys -path, health_url, timeout = sys.argv[1:] +path, health_url, metrics_url, health_timeout, metrics_timeout = sys.argv[1:] wanted = { "SMTP_RELAY", "SMTP_PORT", @@ -59,6 +66,7 @@ wanted = { "EMAIL_FROM_ADDRESS", "EMAIL_FROM_NAME", "SUPPORT_INBOX_ADDRESS", + "METRICS_TOKEN", } values = {} with open(path, encoding="utf-8") as handle: @@ -77,8 +85,11 @@ if missing: raise SystemExit("Missing production mail settings: " + ", ".join(missing)) payload = { - "health_timeout_seconds": int(timeout), + "health_timeout_seconds": int(health_timeout), "health_url": health_url, + "metrics_timeout_seconds": int(metrics_timeout), + "metrics_token": values["METRICS_TOKEN"], + "metrics_url": metrics_url, "smtp": { "from_address": values["EMAIL_FROM_ADDRESS"], "from_name": values["EMAIL_FROM_NAME"], @@ -98,7 +109,7 @@ chmod 600 "$config" cat >"$service" <"$timer" <[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") @@ -51,6 +67,67 @@ def required_text(mapping: dict[str, Any], name: str) -> str: 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): @@ -108,24 +185,55 @@ def check_health(config: dict[str, Any]) -> tuple[str, str]: 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) - status, detail = check_health(config) + 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 readiness is DOWN", + "[Who Need Help] Production monitor is DOWN", "\n".join( [ - "The independent production readiness check failed.", + "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.", @@ -135,22 +243,54 @@ def monitor(config_path: Path, state_path: Path) -> int: elif previous_status == "down": send_message( config, - "[Who Need Help] Production readiness recovered", + "[Who Need Help] Production monitor recovered", "\n".join( [ - "The independent production readiness check recovered.", + "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, }, ) @@ -185,6 +325,7 @@ def send_test_notification(config_path: Path) -> int: [ "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.", ] diff --git a/scripts/quality.sh b/scripts/quality.sh index 027f10d..62d1682 100755 --- a/scripts/quality.sh +++ b/scripts/quality.sh @@ -1529,6 +1529,10 @@ docker run --rm \ --volume "$ROOT/scripts/alert-receiver.py:/src/alert-receiver.py:ro" \ "$PYTHON_IMAGE" python -c \ 'import py_compile; py_compile.compile("/src/alert-receiver.py", cfile="/tmp/alert-receiver.pyc", doraise=True)' +docker run --rm \ + --volume "$ROOT:/src:ro" \ + --workdir /src \ + "$PYTHON_IMAGE" python test/scripts/production_external_monitor_test.py docker run --rm \ --volume "$ROOT/ops/external-boundaries/mock_server.py:/src/mock_server.py:ro" \ "$PYTHON_IMAGE" python -c \ diff --git a/test/scripts/production_external_monitor_test.py b/test/scripts/production_external_monitor_test.py new file mode 100644 index 0000000..0adb3d2 --- /dev/null +++ b/test/scripts/production_external_monitor_test.py @@ -0,0 +1,170 @@ +#!/usr/bin/env python3 + +import importlib.util +import json +import tempfile +import unittest +from pathlib import Path +from unittest import mock + + +SCRIPT = Path(__file__).resolve().parents[2] / "scripts" / "production-external-monitor.py" +SPEC = importlib.util.spec_from_file_location("production_external_monitor", SCRIPT) +assert SPEC is not None and SPEC.loader is not None +MONITOR = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(MONITOR) + + +class ProductionExternalMonitorTest(unittest.TestCase): + def test_parses_only_aggregate_failure_counters(self): + payload = """ +# TYPE who_need_help_http_exceptions_total counter +who_need_help_http_exceptions_total 2 +who_need_help_oban_jobs_failed_total{queue="push"} 3 +who_need_help_oban_jobs_failed_total{queue="mailers"} 4 +who_need_help_email_deliveries_total{status="ok"} 50 +who_need_help_email_deliveries_total{status="error"} 5 +who_need_help_email_delivery_exceptions_total 1 +who_need_help_http_requests_total 999 +""" + + self.assertEqual( + MONITOR.parse_monitored_counters(payload), + { + "who_need_help_http_exceptions_total": 2.0, + "who_need_help_oban_jobs_failed_total|queue=push": 3.0, + "who_need_help_oban_jobs_failed_total|queue=mailers": 4.0, + "who_need_help_email_deliveries_total|status=error": 5.0, + "who_need_help_email_delivery_exceptions_total": 1.0, + }, + ) + + def test_first_observation_establishes_a_baseline(self): + self.assertEqual( + MONITOR.counter_increases(None, {"failure": 4.0}), + {}, + ) + + def test_new_and_increased_failures_are_reported_after_baseline(self): + self.assertEqual( + MONITOR.counter_increases( + {"existing": 2.0}, + {"existing": 5.0, "new": 1.0}, + ), + {"existing": 3.0, "new": 1.0}, + ) + + def test_counter_reset_does_not_create_a_false_failure(self): + self.assertEqual( + MONITOR.counter_increases({"failure": 10.0}, {"failure": 1.0}), + {}, + ) + + def test_malformed_labels_are_rejected(self): + with self.assertRaises(ValueError): + MONITOR.parse_monitored_counters( + 'who_need_help_oban_jobs_failed_total{queue="push" broken} 1' + ) + + def test_monitor_baselines_then_notifies_only_for_new_failures(self): + with tempfile.TemporaryDirectory() as directory: + config_path = Path(directory) / "config.json" + state_path = Path(directory) / "state.json" + config_path.write_text(json.dumps({}), encoding="utf-8") + messages = [] + + with ( + mock.patch.object(MONITOR, "check_health", return_value=("up", "ready")), + mock.patch.object( + MONITOR, + "check_metrics", + return_value=("up", "parsed", "node-a", {"failure": 2.0}), + ), + mock.patch.object( + MONITOR, + "send_message", + side_effect=lambda _config, subject, body: messages.append((subject, body)), + ), + mock.patch("builtins.print"), + ): + self.assertEqual(MONITOR.monitor(config_path, state_path), 0) + + self.assertEqual(messages, []) + + with ( + mock.patch.object(MONITOR, "check_health", return_value=("up", "ready")), + mock.patch.object( + MONITOR, + "check_metrics", + return_value=("up", "parsed", "node-a", {"failure": 5.0}), + ), + mock.patch.object( + MONITOR, + "send_message", + side_effect=lambda _config, subject, body: messages.append((subject, body)), + ), + mock.patch("builtins.print"), + ): + self.assertEqual(MONITOR.monitor(config_path, state_path), 0) + + self.assertEqual(len(messages), 1) + self.assertIn("failure counters increased", messages[0][0]) + self.assertIn("failure: +3", messages[0][1]) + + def test_monitor_alerts_once_for_down_state_and_once_for_recovery(self): + with tempfile.TemporaryDirectory() as directory: + config_path = Path(directory) / "config.json" + state_path = Path(directory) / "state.json" + config_path.write_text( + json.dumps( + { + "health_url": "https://example.test/healthz/ready", + "metrics_url": "https://example.test/metrics", + } + ), + encoding="utf-8", + ) + messages = [] + + with ( + mock.patch.object(MONITOR, "check_health", return_value=("up", "ready")), + mock.patch.object( + MONITOR, + "check_metrics", + return_value=("down", "HTTP 401", "unknown", {}), + ), + mock.patch.object( + MONITOR, + "send_message", + side_effect=lambda _config, subject, body: messages.append((subject, body)), + ), + mock.patch("builtins.print"), + ): + self.assertEqual(MONITOR.monitor(config_path, state_path), 1) + self.assertEqual(MONITOR.monitor(config_path, state_path), 1) + + self.assertEqual(len(messages), 1) + self.assertIn("DOWN", messages[0][0]) + + with ( + mock.patch.object(MONITOR, "check_health", return_value=("up", "ready")), + mock.patch.object( + MONITOR, + "check_metrics", + return_value=("up", "parsed", "node-a", {}), + ), + mock.patch.object( + MONITOR, + "send_message", + side_effect=lambda _config, subject, body: messages.append((subject, body)), + ), + mock.patch("builtins.print"), + ): + self.assertEqual(MONITOR.monitor(config_path, state_path), 0) + + self.assertEqual(len(messages), 2) + self.assertIn("recovered", messages[1][0]) + + +if __name__ == "__main__": + unittest.main() diff --git a/test/who_need_help_web/controllers/metrics_controller_test.exs b/test/who_need_help_web/controllers/metrics_controller_test.exs index 7b124b4..25639e6 100644 --- a/test/who_need_help_web/controllers/metrics_controller_test.exs +++ b/test/who_need_help_web/controllers/metrics_controller_test.exs @@ -43,6 +43,50 @@ defmodule WhoNeedHelpWeb.MetricsControllerTest do assert body =~ "who_need_help_database_query_decode_duration_microseconds_total " end + test "exports aggregate email delivery outcomes without recipient labels" do + reporter_name = :email_delivery_metrics_controller_test + + metrics = + WhoNeedHelpWeb.Telemetry.prometheus_metrics() + |> Enum.filter(fn metric -> + metric.name in [ + [:who_need_help, :email, :deliveries, :total], + [:who_need_help, :email, :delivery, :exceptions, :total] + ] + end) + + start_supervised!( + {TelemetryMetricsPrometheus.Core, name: reporter_name, metrics: metrics, start_async: false} + ) + + :telemetry.execute( + [:swoosh, :deliver, :stop], + %{duration: 100}, + %{mailer: WhoNeedHelp.Mailer, result: %{id: "accepted"}} + ) + + :telemetry.execute( + [:swoosh, :deliver, :stop], + %{duration: 100}, + %{mailer: WhoNeedHelp.Mailer, error: :rejected} + ) + + :telemetry.execute( + [:swoosh, :deliver, :exception], + %{duration: 100}, + %{mailer: WhoNeedHelp.Mailer, kind: :error, reason: :timeout} + ) + + body = TelemetryMetricsPrometheus.Core.scrape(reporter_name) + + assert body =~ ~s(who_need_help_email_deliveries_total{status="ok"} 1) + assert body =~ ~s(who_need_help_email_deliveries_total{status="error"} 1) + assert body =~ "who_need_help_email_delivery_exceptions_total 1" + refute body =~ "accepted" + refute body =~ "rejected" + refute body =~ "recipient" + end + test "database execution metric tolerates events without query_time" do metric = WhoNeedHelpWeb.Telemetry.prometheus_metrics()