Monitor aggregate production delivery failures

This commit is contained in:
SimpleTest 2026-08-09 07:03:45 +03:00
parent a65d7d126f
commit e779188e8e
7 changed files with 433 additions and 29 deletions

View File

@ -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

View File

@ -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 ->

View File

@ -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" <<EOF
[Unit]
Description=Who Need Help independent production readiness monitor
Description=Who Need Help independent production operations monitor
Wants=network-online.target
After=network-online.target
@ -109,7 +120,7 @@ EOF
cat >"$timer" <<EOF
[Unit]
Description=Run the Who Need Help independent production readiness monitor
Description=Run the Who Need Help independent production operations monitor
[Timer]
OnCalendar=$monitor_calendar
@ -133,5 +144,7 @@ ssh -o BatchMode=yes "$monitor_target" \
printf 'External monitor installed on %s (%s).\n' "$monitor_target" "$monitor_host"
printf 'Health URL: %s\n' "$monitor_url"
printf 'Schedule: %s; request timeout: %ss.\n' "$monitor_calendar" "$health_timeout"
echo "The SMTP credential is stored only in a mode-0600 configuration on the external monitor host."
printf 'Metrics URL: %s\n' "$metrics_url"
printf 'Schedule: %s; health timeout: %ss; metrics timeout: %ss.\n' \
"$monitor_calendar" "$health_timeout" "$metrics_timeout"
echo "The SMTP and metrics credentials are stored only in a mode-0600 configuration on the external monitor host."

View File

@ -6,6 +6,7 @@ from __future__ import annotations
import argparse
import json
import os
import re
import smtplib
import ssl
import sys
@ -18,6 +19,21 @@ 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")
@ -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.",
]

View File

@ -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 \

View File

@ -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()

View File

@ -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()