"""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} # --- product events ----------------------------------------------------------- # These count intent (a tap, a signup), not requests, and are keyed by # (event, referrer domain, UTC day). The endpoint is unauthenticated, so the # allowlist and the caps are load-bearing, not decoration. def test_event_allowlist(): m = _fresh() assert m.normalize_event("home.locate") == "home.locate" assert m.normalize_event("home.nav_calendar") == "home.nav_calendar" # Unknown names collapse into one bucket rather than growing the key space. assert m.normalize_event("home.nav_not_a_surface") == m.EVENT_OTHER assert m.normalize_event("totally.made.up") == m.EVENT_OTHER assert m.normalize_event("") == m.EVENT_OTHER assert m.normalize_event(None) == m.EVENT_OTHER def test_referrer_normalization(): m = _fresh() # Only the bare domain is kept — never the path or query, which can carry a # search term or a private URL. assert m.normalize_referrer("https://news.ycombinator.com/item?id=1") == "news.ycombinator.com" assert m.normalize_referrer("https://www.reddit.com/r/weather/") == "reddit.com" assert m.normalize_referrer(None) == "direct" assert m.normalize_referrer("") == "direct" assert m.normalize_referrer("https://thermograph.org/x", host="thermograph.org") == "self" assert m.normalize_referrer("https://www.thermograph.org/x", host="thermograph.org") == "self" # Junk is sanitized, never stored raw. assert " " not in m.normalize_referrer("https://ev il<>.com/") assert len(m.normalize_referrer("https://" + "a" * 200 + ".com/")) <= 64 def test_events_are_not_counted_as_inbound_traffic(): """The beacon reports traffic; counting it would double every interaction.""" m = _fresh() assert m.classify_inbound("/thermograph/api/v2/event", "/thermograph") == "event" m.record_inbound("event", 204) assert m.snapshot()["inbound_total"] == 0 assert "event" not in m.snapshot()["inbound"] def test_record_event_keys_by_event_referrer_day(): m = _fresh() for _ in range(3): m.record_event("home.locate", referer="https://news.ycombinator.com/x") m.record_event("home.share", referer=None) events = m.snapshot()["events"] day = m._utc_day() assert events["home.locate"]["news.ycombinator.com"][day] == 3 assert events["home.share"]["direct"][day] == 1 assert m.snapshot()["events_total"] == 4 def test_unknown_event_carries_no_referrer_dimension(): m = _fresh() m.record_event("junk.event", referer="https://spam.example/") events = m.snapshot()["events"] assert list(events[m.EVENT_OTHER]) == ["direct"] def test_referrer_cardinality_is_capped(): """A spray of forged Referer headers must not grow the key space forever.""" m = _fresh() m._EVENT_RATE_PER_MIN = 10_000 for i in range(m.MAX_REFERRERS + 50): m.record_event("home.locate", referer=f"https://ref{i}.example/") refs = m.snapshot()["events"]["home.locate"] assert len(refs) <= m.MAX_REFERRERS + 1 # +1 for the overflow bucket assert m.REFERRER_OTHER in refs def test_event_rate_limit_per_ip(): m = _fresh() for _ in range(m._EVENT_RATE_PER_MIN + 25): m.record_event("home.locate", ip="203.0.113.9") assert m.snapshot()["events_total"] == m._EVENT_RATE_PER_MIN