The ops dashboard marked subscription-notifier DOWN ~2/3 of the time after the notifier was gated to a single leader worker. It judged the daemon by whether its thread was present on whichever worker answered /api/v2/metrics, but only the leader runs it — so polls landing on the other two workers saw no thread and reported DOWN for a perfectly healthy notifier. Per-worker thread checks can't describe a single-leader daemon, and worse, they could no longer distinguish a real crash from a normal non-leader. Report it by liveness instead: the notifier stamps a heartbeat (timestamp + expected interval) into the shared metrics store each loop iteration, including once at startup so the dashboard shows it up immediately after a restart. Since the beat lives in the shared SQLite store, any worker's metrics response carries it, so the dashboard sees the notifier's health whichever worker it polls. - metrics.py: both stores record/expose heartbeats (SQLite reuses the meta table, so every worker's beats share one file); snapshot() converts each beat to a server-computed age so readers needn't reconcile clock skew. - notify.py: beat at loop start and every wake. - dashboard.py: judge heartbeat daemons by freshness (up to 2 missed intervals before DOWN, with beat age shown); neighbor-warmer keeps its per-worker thread check, which is correct since every worker runs it. - Tests: heartbeat store/snapshot shape + cross-worker sharing; dashboard fresh/stale/missing rendering and the unchanged warmer check. 194 passed. Co-authored-by: root <root@vmi3417050.contaboserver.net>
241 lines
10 KiB
Python
241 lines
10 KiB
Python
"""Unit tests for the in-process traffic counters (metrics.py)."""
|
|
import importlib
|
|
|
|
import metrics
|
|
|
|
|
|
def _fresh():
|
|
# Counters are module-level; reload for an isolated starting state per test.
|
|
return importlib.reload(metrics)
|
|
|
|
|
|
def test_phase_source_map_covers_every_source():
|
|
m = _fresh()
|
|
assert m.source_for_phase("history_fetch") == "open-meteo-archive"
|
|
assert m.source_for_phase("history_topup") == "open-meteo-archive"
|
|
assert m.source_for_phase("recent_forecast_fetch") == "open-meteo-forecast"
|
|
assert m.source_for_phase("geocode") == "open-meteo-geocoding"
|
|
assert m.source_for_phase("history_nasa") == "nasa-power"
|
|
assert m.source_for_phase("forecast_metno") == "met-norway"
|
|
assert m.source_for_phase("reverse_geocode") == "nominatim"
|
|
assert m.source_for_phase("unknown-phase") == "other"
|
|
|
|
|
|
def test_classify_inbound_categories():
|
|
m = _fresh()
|
|
c = m.classify_inbound
|
|
# BASE stripped, API namespaces
|
|
assert c("/thermograph/api/v2/place", "/thermograph") == "api:place"
|
|
assert c("/thermograph/api/v2/grade", "/thermograph") == "api:grade"
|
|
assert c("/thermograph/api/v1/geocode", "/thermograph") == "api:geocode"
|
|
assert c("/thermograph/api/v2/auth/login", "/thermograph") == "auth"
|
|
assert c("/thermograph/api/v2/users/me", "/thermograph") == "auth"
|
|
assert c("/thermograph/api/v2/subscriptions", "/thermograph") == "accounts"
|
|
assert c("/thermograph/api/v2/notifications", "/thermograph") == "accounts"
|
|
assert c("/thermograph/api/v2/push/subscribe", "/thermograph") == "accounts"
|
|
assert c("/thermograph/api/v2/metrics", "/thermograph") == "metrics"
|
|
# The metrics endpoint is never counted, even under an unexpected prefix — e.g. the
|
|
# dashboard probing /thermograph/... against a root-served (base="") prod app.
|
|
assert c("/thermograph/api/v2/metrics", "") == "metrics"
|
|
assert c("/api/v2/metrics", "") == "metrics"
|
|
# pages, static, seo
|
|
assert c("/thermograph/", "/thermograph") == "page"
|
|
assert c("/thermograph/calendar", "/thermograph") == "page"
|
|
assert c("/thermograph/climate/tokyo", "/thermograph") == "page"
|
|
assert c("/thermograph/app.js", "/thermograph") == "static"
|
|
assert c("/thermograph/style.css", "/thermograph") == "static"
|
|
assert c("/thermograph/robots.txt", "/thermograph") == "seo"
|
|
# root-mounted (BASE == "")
|
|
assert c("/api/v2/place", "") == "api:place"
|
|
assert c("/", "") == "page"
|
|
|
|
|
|
def test_record_inbound_buckets_by_status():
|
|
m = _fresh()
|
|
m.record_inbound("api:place", 200)
|
|
m.record_inbound("api:place", 200)
|
|
m.record_inbound("api:place", 503)
|
|
m.record_inbound("page", 304)
|
|
snap = m.snapshot()
|
|
assert snap["inbound"]["api:place"]["total"] == 3
|
|
assert snap["inbound"]["api:place"]["2xx"] == 2
|
|
assert snap["inbound"]["api:place"]["5xx"] == 1
|
|
assert snap["inbound"]["page"]["3xx"] == 1
|
|
assert snap["inbound_total"] == 4
|
|
|
|
|
|
def test_metrics_endpoint_polling_is_not_counted():
|
|
# The dashboard polls the metrics endpoint; that self-traffic must be excluded.
|
|
m = _fresh()
|
|
m.record_inbound("metrics", 200)
|
|
m.record_inbound("metrics", 404)
|
|
m.record_inbound("api:place", 200)
|
|
snap = m.snapshot()
|
|
assert "metrics" not in snap["inbound"]
|
|
assert snap["inbound_total"] == 1
|
|
|
|
|
|
def test_record_outbound_by_source_and_outcome():
|
|
m = _fresh()
|
|
m.record_outbound("history_fetch", "ok")
|
|
m.record_outbound("history_topup", "ok") # same source (archive)
|
|
m.record_outbound("reverse_geocode", "error")
|
|
m.record_outbound("recent_forecast_fetch", "rate_limited")
|
|
snap = m.snapshot()
|
|
assert snap["outbound"]["open-meteo-archive"]["ok"] == 2
|
|
assert snap["outbound"]["nominatim"]["error"] == 1
|
|
assert snap["outbound"]["open-meteo-forecast"]["rate_limited"] == 1
|
|
assert snap["outbound_total"] == 4
|
|
|
|
|
|
def test_snapshot_has_process_facts():
|
|
m = _fresh()
|
|
snap = m.snapshot()
|
|
assert isinstance(snap["pid"], int)
|
|
assert snap["uptime_s"] >= 0
|
|
assert "open-meteo-archive" in snap["outbound"] # stable source rows exist
|
|
|
|
|
|
def test_recent_window_sums_and_rolls_off(monkeypatch):
|
|
m = _fresh()
|
|
clock = {"now": 1_000_000.0}
|
|
monkeypatch.setattr(m.time, "time", lambda: clock["now"])
|
|
|
|
# minute 0: three page hits + one archive call
|
|
for _ in range(3):
|
|
m.record_inbound("page", 200)
|
|
m.record_outbound("history_fetch", "ok") # -> open-meteo-archive
|
|
snap = m.snapshot()
|
|
assert snap["window_minutes"] == 10
|
|
assert snap["inbound_recent"]["page"] == 3
|
|
assert snap["outbound_recent"]["open-meteo-archive"] == 1
|
|
|
|
# +5 min: still inside the 10-minute window, so the earlier hits still count
|
|
clock["now"] += 5 * 60
|
|
m.record_inbound("page", 200)
|
|
assert m.snapshot()["inbound_recent"]["page"] == 4
|
|
|
|
# +6 more min (t = 11 min): the minute-0 burst has rolled out of the window;
|
|
# only the hit from minute 5 remains, and the archive call is gone entirely.
|
|
clock["now"] += 6 * 60
|
|
snap = m.snapshot()
|
|
assert snap["inbound_recent"].get("page", 0) == 1
|
|
assert "open-meteo-archive" not in snap["outbound_recent"]
|
|
# the cumulative since-start total is unaffected by the rolling window
|
|
assert snap["inbound"]["page"]["total"] == 4
|
|
|
|
|
|
def test_recent_excludes_self_polling():
|
|
m = _fresh()
|
|
m.record_inbound("metrics", 200) # dashboard polling — never counted
|
|
m.record_inbound("page", 200)
|
|
snap = m.snapshot()
|
|
assert "metrics" not in snap["inbound_recent"]
|
|
assert snap["inbound_recent"]["page"] == 1
|
|
|
|
|
|
# --- shared SQLite store (multi-worker) -------------------------------------------
|
|
# When the app runs several uvicorn workers, per-process memory would split the tally
|
|
# across workers; THERMOGRAPH_METRICS_DB points them at one shared SQLite file instead.
|
|
|
|
def test_sqlite_store_matches_snapshot_shape(tmp_path):
|
|
m = _fresh()
|
|
store = m._SqliteStore(str(tmp_path / "metrics.db"))
|
|
store.record_inbound("api:place", "2xx")
|
|
store.record_inbound("api:place", "2xx")
|
|
store.record_inbound("api:place", "5xx")
|
|
store.record_outbound("open-meteo-archive", "ok")
|
|
data = store.read()
|
|
assert data["inbound"]["api:place"]["total"] == 3
|
|
assert data["inbound"]["api:place"]["2xx"] == 2
|
|
assert data["inbound"]["api:place"]["5xx"] == 1
|
|
# stable source rows are present even at zero traffic, like the in-memory store
|
|
assert data["outbound"]["open-meteo-archive"]["ok"] == 1
|
|
assert "nasa-power" in data["outbound"]
|
|
assert data["inbound_recent"]["api:place"] == 3
|
|
assert data["outbound_recent"]["open-meteo-archive"] == 1
|
|
|
|
|
|
def test_sqlite_store_is_shared_across_workers(tmp_path):
|
|
# Two store instances on the same file stand in for two worker processes: what one
|
|
# records, the other must see — that's the whole point of the shared store.
|
|
m = _fresh()
|
|
db = str(tmp_path / "metrics.db")
|
|
worker_a = m._SqliteStore(db)
|
|
worker_b = m._SqliteStore(db)
|
|
worker_a.record_inbound("page", "2xx")
|
|
worker_b.record_inbound("page", "2xx")
|
|
worker_b.record_inbound("api:grade", "2xx")
|
|
# Either worker, polled by the dashboard, reports the combined tally.
|
|
for worker in (worker_a, worker_b):
|
|
data = worker.read()
|
|
assert data["inbound"]["page"]["total"] == 2
|
|
assert data["inbound"]["api:grade"]["total"] == 1
|
|
assert data["inbound_recent"]["page"] == 2
|
|
# Both workers share one start time (first writer wins), not one-per-process.
|
|
assert worker_a.started_at == worker_b.started_at
|
|
|
|
|
|
def test_sqlite_recent_window_rolls_off(tmp_path, monkeypatch):
|
|
m = _fresh()
|
|
clock = {"now": 1_000_000.0}
|
|
monkeypatch.setattr(m.time, "time", lambda: clock["now"])
|
|
store = m._SqliteStore(str(tmp_path / "metrics.db"))
|
|
|
|
store.record_inbound("page", "2xx")
|
|
assert store.read()["inbound_recent"]["page"] == 1
|
|
# +11 min: the hit has rolled out of the 10-minute window and is pruned...
|
|
clock["now"] += 11 * 60
|
|
store.record_inbound("page", "2xx") # advances the minute -> triggers prune
|
|
data = store.read()
|
|
assert data["inbound_recent"]["page"] == 1 # only the fresh hit remains in-window
|
|
assert data["inbound"]["page"]["total"] == 2 # ...but the cumulative total is intact
|
|
|
|
|
|
def test_env_selects_sqlite_store(tmp_path, monkeypatch):
|
|
# Setting the env var makes the module boot on the shared store, not in-memory.
|
|
monkeypatch.setenv("THERMOGRAPH_METRICS_DB", str(tmp_path / "metrics.db"))
|
|
m = _fresh()
|
|
assert isinstance(m._store, m._SqliteStore)
|
|
m.record_inbound("page", 200)
|
|
assert m.snapshot()["inbound"]["page"]["total"] == 1
|
|
|
|
|
|
def test_heartbeat_snapshot_reports_age_and_interval(monkeypatch):
|
|
# A daemon stamps its liveness; the snapshot reports it as an age on the server clock
|
|
# (so a reader needn't reconcile clock skew) plus the beat's expected interval.
|
|
m = _fresh()
|
|
clock = {"now": 1_000_000.0}
|
|
monkeypatch.setattr(m.time, "time", lambda: clock["now"])
|
|
m.record_heartbeat("subscription-notifier", 900)
|
|
clock["now"] += 42 # 42 s pass before the dashboard polls
|
|
hb = m.snapshot()["heartbeats"]["subscription-notifier"]
|
|
assert hb["age_s"] == 42.0
|
|
assert hb["interval_s"] == 900
|
|
|
|
|
|
def test_heartbeat_absent_until_recorded():
|
|
# No beat yet => no entry (the dashboard renders this as DOWN / "no heartbeat").
|
|
m = _fresh()
|
|
assert m.snapshot()["heartbeats"] == {}
|
|
|
|
|
|
def test_heartbeat_shared_across_workers_sqlite(tmp_path):
|
|
# Only the leader worker beats, but every worker reads the shared store — so the
|
|
# dashboard sees the notifier alive whichever worker happens to answer its poll.
|
|
m = _fresh()
|
|
db = str(tmp_path / "metrics.db")
|
|
leader = m._SqliteStore(db)
|
|
non_leader = m._SqliteStore(db) # never beats (it isn't the notifier leader)
|
|
leader.record_heartbeat("subscription-notifier", 5_000.0, 900)
|
|
beats = non_leader.read()["heartbeats"]
|
|
assert beats["subscription-notifier"]["ts"] == 5_000.0
|
|
assert beats["subscription-notifier"]["interval_s"] == 900
|
|
|
|
|
|
def test_heartbeat_survives_store_read_shape(tmp_path):
|
|
# The in-memory and SQLite stores expose heartbeats identically (ts + interval_s).
|
|
m = _fresh()
|
|
mem = m._MemStore()
|
|
mem.record_heartbeat("subscription-notifier", 10.0, 900)
|
|
assert mem.read()["heartbeats"]["subscription-notifier"] == {"ts": 10.0, "interval_s": 900}
|