Extends the existing /api/v2/event beacon into a small, typed, durable event schema so the interactive views can answer product questions the request log cannot: whether a visitor ever gets a location, whether search finds anything, which controls and views earn their maintenance, and which dead ends people actually hit. Seven events with three enum dimension slots, stored as hourly aggregates in a TimescaleDB hypertable plus one JSONL line per event for the 30-day Loki view. No row per interaction and no column an identifier could go in: no cookie, no session id, no user id, no IP, no coordinates, no free text, no URLs. The IP is a per-minute rate-limit key and nothing else. Abuse resistance for a public endpoint on a 60%-crawler site: closed allowlist (unknown names collapse to one bucket), 4 KB streaming body cap, per-IP ceiling raised to 120/min (30 was sized for four coarse events and would have silently truncated a batched stream), soft Sec-Fetch-Site first-party check, uniform 204. Ships inert -- THERMOGRAPH_EVENTS gates both the recorder and the flag stamp on <html> that makes the client send anything, and is unset everywhere. See UI-EVENTS.md for the schema, the rejected transport/storage alternatives, and the privacy decisions that must be settled before the flag is turned on.
822 lines
38 KiB
Python
822 lines
38 KiB
Python
"""Traffic counters for the ops dashboard.
|
|
|
|
Cheap, best-effort counters for **inbound** HTTP requests (bucketed by category) and
|
|
**outbound** calls to external data sources (bucketed by source). Everything here is
|
|
since-start; durable daily history lives in the audit/error JSONL logs (see
|
|
``audit.py``). Recording never raises into the request path.
|
|
|
|
Read over the gated ``GET {BASE}/api/v2/metrics`` route (see ``app.py``) — a raw
|
|
counters/heartbeats API. (Dashboards proper now live in the separate
|
|
``thermograph-observability`` project, a Grafana + Loki stack fed from the JSON
|
|
logs; this endpoint remains for quick token-gated checks over an SSH tunnel.)
|
|
|
|
**Multi-worker.** Under a single uvicorn worker the counters live in process memory.
|
|
When the app runs several workers (``--workers N``), per-process memory would split the
|
|
tally across workers, so the dashboard — which polls one random worker — would see only
|
|
a fraction of the traffic. Set ``THERMOGRAPH_METRICS_DB`` to a file path and every worker
|
|
shares one SQLite store, so the dashboard reports the whole picture. Unset (dev, tests,
|
|
single worker) keeps the zero-dependency in-memory store.
|
|
"""
|
|
import os
|
|
import sqlite3
|
|
import threading
|
|
import time
|
|
import urllib.parse
|
|
|
|
# Dialect selection mirrors accounts/db.py and data/store.py: a ``postgresql://``
|
|
# URL in THERMOGRAPH_DATABASE_URL routes the shared counters into Postgres (prod /
|
|
# multi-worker); otherwise the store stays SQLite-file / in-memory as before.
|
|
DATABASE_URL = os.environ.get("THERMOGRAPH_DATABASE_URL", "").strip()
|
|
IS_POSTGRES = DATABASE_URL.startswith("postgresql")
|
|
|
|
|
|
def _libpq_dsn(url: str) -> str:
|
|
"""A plain libpq URL that ``psycopg.connect`` accepts: drop the SQLAlchemy
|
|
driver suffix (``postgresql+asyncpg://`` / ``+psycopg://`` -> ``postgresql://``)."""
|
|
return url.replace("+asyncpg", "").replace("+psycopg", "")
|
|
|
|
|
|
_PG_DSN = _libpq_dsn(DATABASE_URL) if IS_POSTGRES else ""
|
|
|
|
# The ``phase=`` tag passed to ``climate._request`` maps 1:1 onto an external source.
|
|
# Keep this in sync with the phase names used at each call site in ``climate.py``.
|
|
_PHASE_SOURCE = {
|
|
"history_fetch": "open-meteo-archive",
|
|
"history_topup": "open-meteo-archive",
|
|
"recent_forecast_fetch": "open-meteo-forecast",
|
|
# Forward geocoding moved off Open-Meteo to Nominatim (same host as reverse) —
|
|
# /geocode now backstops the local GeoNames index on a miss.
|
|
"geocode": "nominatim",
|
|
"history_nasa": "nasa-power",
|
|
"forecast_metno": "met-norway",
|
|
# Wind gusts for the gust-less backup sources (NASA POWER / MET Norway), read
|
|
# from Meteostat's keyless bulk endpoints (station list + per-station daily).
|
|
"history_gust": "meteostat",
|
|
"meteostat_stations": "meteostat",
|
|
"reverse_geocode": "nominatim",
|
|
}
|
|
|
|
# Every source we know about, so the dashboard shows a stable set of rows even before
|
|
# any traffic has hit a given source.
|
|
SOURCES = sorted(set(_PHASE_SOURCE.values()))
|
|
|
|
_OUTCOMES = ("ok", "retry", "error", "rate_limited")
|
|
_STATUS_BUCKETS = ("2xx", "3xx", "4xx", "5xx", "other")
|
|
|
|
# Rolling short-window view: per-minute buckets summed over the last WINDOW minutes on
|
|
# snapshot(). Cheap and bounded — minutes are pruned as they roll out of the window.
|
|
WINDOW_MINUTES = 10
|
|
|
|
|
|
def source_for_phase(phase: str) -> str:
|
|
"""Map a ``climate._request`` phase tag to an external-source label."""
|
|
return _PHASE_SOURCE.get(phase, "other")
|
|
|
|
|
|
def _status_bucket(status) -> str:
|
|
try:
|
|
code = int(status)
|
|
except (TypeError, ValueError):
|
|
return "other"
|
|
if 200 <= code < 300:
|
|
return "2xx"
|
|
if 300 <= code < 400:
|
|
return "3xx"
|
|
if 400 <= code < 500:
|
|
return "4xx"
|
|
if 500 <= code < 600:
|
|
return "5xx"
|
|
return "other"
|
|
|
|
|
|
def classify_inbound(path: str, base: str = "") -> str:
|
|
"""Bucket a request path into a dashboard category.
|
|
|
|
``path`` is ``request.url.path`` (still carrying the ``BASE`` prefix); ``base`` is
|
|
the app's mount prefix (e.g. ``/thermograph`` or ``""``).
|
|
"""
|
|
p = path or "/"
|
|
# The dashboard polls the metrics endpoint; never count it as traffic, whatever base
|
|
# prefix it arrives under — e.g. a `/thermograph/api/v2/metrics` probe against a
|
|
# root-served prod app would otherwise land in "other" and show as an inbound error.
|
|
if p.rstrip("/").endswith("/api/v2/metrics"):
|
|
return "metrics"
|
|
if base and p.startswith(base):
|
|
p = p[len(base):] or "/"
|
|
if p.startswith("/api/"):
|
|
rest = p[len("/api/"):]
|
|
if rest.startswith(("v1/", "v2/")):
|
|
rest = rest[3:]
|
|
seg = rest.split("/", 1)[0].split("?", 1)[0]
|
|
if seg in ("auth", "users"):
|
|
return "auth"
|
|
if seg in ("subscriptions", "notifications", "push"):
|
|
return "accounts"
|
|
if seg == "metrics":
|
|
return "metrics"
|
|
# The event beacon reports traffic; it is not itself traffic. Its own
|
|
# category keeps record_inbound from double-counting every interaction.
|
|
if seg == "event":
|
|
return "event"
|
|
return f"api:{seg}" if seg else "api:other"
|
|
if p.endswith((".js", ".css", ".html", ".webmanifest", ".png", ".svg", ".ico",
|
|
".json", ".woff", ".woff2", ".map", ".txt", ".xml")):
|
|
# robots.txt / sitemap.xml are crawler SEO endpoints, not page views.
|
|
if p in ("/robots.txt", "/sitemap.xml"):
|
|
return "seo"
|
|
return "static"
|
|
if p in ("", "/", "/calendar", "/day", "/compare", "/legend", "/alerts", "/privacy"):
|
|
return "page"
|
|
if p == "/about" or p.startswith("/climate") or p.startswith("/glossary"):
|
|
return "page"
|
|
return "other"
|
|
|
|
|
|
# --- product events ----------------------------------------------------------
|
|
# Distinct from inbound traffic: these count *intent* (tapped "use my location",
|
|
# submitted the digest form) rather than requests. Keyed by (event, referrer
|
|
# domain, UTC day) so a launch post's traffic can be told apart from search.
|
|
#
|
|
# Everything here is aggregate — no cookies, no per-visitor id, no full URLs.
|
|
|
|
EVENTS = frozenset({
|
|
"home.locate", "home.digest_signup", "home.records_click", "home.share",
|
|
})
|
|
# home.nav_{surface}: only these surfaces, so the key space stays bounded.
|
|
NAV_SURFACES = frozenset({
|
|
"calendar", "compare", "day", "alerts", "city", "cities", "records",
|
|
})
|
|
# Anything unrecognized collapses here, with no referrer dimension — visible
|
|
# enough to notice abuse, bounded enough to be harmless.
|
|
EVENT_OTHER = "event:other"
|
|
|
|
EVENT_DAYS = 7 # how many days of event history to keep
|
|
MAX_REFERRERS = 200 # distinct referrer domains tracked before overflow
|
|
REFERRER_OTHER = "other"
|
|
|
|
_REF_OK = set("abcdefghijklmnopqrstuvwxyz0123456789.-")
|
|
|
|
|
|
def normalize_event(name: str) -> str:
|
|
"""An allowlisted event name, or EVENT_OTHER. The allowlist is what keeps an
|
|
unauthenticated endpoint from being turned into unbounded storage."""
|
|
name = (name or "").strip()
|
|
if name in EVENTS:
|
|
return name
|
|
surface = name.partition("home.nav_")[2] if name.startswith("home.nav_") else ""
|
|
if surface and surface in NAV_SURFACES:
|
|
return name
|
|
return EVENT_OTHER
|
|
|
|
|
|
def normalize_referrer(referer: "str | None", host: "str | None" = None) -> str:
|
|
"""The bare registrable-ish domain of a Referer header — never the full URL,
|
|
which can carry a search query or a private path. 'direct' when absent,
|
|
'self' for our own pages."""
|
|
if not referer:
|
|
return "direct"
|
|
try:
|
|
netloc = urllib.parse.urlsplit(referer).netloc
|
|
except ValueError:
|
|
return REFERRER_OTHER
|
|
domain = (netloc.rpartition("@")[2].partition(":")[0] or "").lower()
|
|
if domain.startswith("www."):
|
|
domain = domain[4:]
|
|
if host and domain == (host.partition(":")[0] or "").lower().removeprefix("www."):
|
|
return "self"
|
|
domain = "".join(ch for ch in domain if ch in _REF_OK)[:64]
|
|
return domain or REFERRER_OTHER
|
|
|
|
|
|
def _utc_day() -> str:
|
|
return time.strftime("%Y-%m-%d", time.gmtime())
|
|
|
|
|
|
def _recent_days() -> "set[str]":
|
|
"""The EVENT_DAYS UTC days still worth keeping."""
|
|
now = time.time()
|
|
return {time.strftime("%Y-%m-%d", time.gmtime(now - d * 86400))
|
|
for d in range(EVENT_DAYS)}
|
|
|
|
|
|
def _nest_events(rows) -> dict:
|
|
"""Flat (event, referrer, day, count) rows -> {event: {referrer: {day: n}}}."""
|
|
out: "dict[str, dict[str, dict[str, int]]]" = {}
|
|
for key, count in rows:
|
|
event, referrer, day = key
|
|
out.setdefault(event, {}).setdefault(referrer, {})[day] = count
|
|
return out
|
|
|
|
|
|
def _now_minute() -> int:
|
|
"""Current epoch-minute. Uses ``time.time()`` so tests can monkeypatch the clock."""
|
|
return int(time.time() // 60)
|
|
|
|
|
|
def _empty_inbound_row() -> "dict[str, int]":
|
|
return {"total": 0, **{b: 0 for b in _STATUS_BUCKETS}}
|
|
|
|
|
|
class _MemStore:
|
|
"""In-process counters. One per worker — fine for a single-worker deployment."""
|
|
|
|
def __init__(self) -> None:
|
|
self._lock = threading.Lock()
|
|
self.started_at = time.time()
|
|
self._inbound: "dict[str, dict[str, int]]" = {}
|
|
self._outbound: "dict[str, dict[str, int]]" = {
|
|
s: {o: 0 for o in _OUTCOMES} for s in SOURCES
|
|
}
|
|
# epoch-minute -> {key: count}, pruned as minutes roll out of the window.
|
|
self._recent_in: "dict[int, dict[str, int]]" = {}
|
|
self._recent_out: "dict[int, dict[str, int]]" = {}
|
|
# name -> {"ts": last-beat epoch, "interval_s": expected cadence}. Lets the
|
|
# dashboard judge a background daemon by liveness, not per-worker thread presence.
|
|
self._heartbeats: "dict[str, dict[str, float]]" = {}
|
|
# (event, referrer, day) -> count, pruned to EVENT_DAYS on write.
|
|
self._events: "dict[tuple[str, str, str], int]" = {}
|
|
self._recent_ev: "dict[int, dict[str, int]]" = {}
|
|
|
|
def _bump_recent(self, buckets, key) -> None:
|
|
"""Count one event for ``key`` in the current-minute bucket (caller holds lock)."""
|
|
minute = _now_minute()
|
|
row = buckets.get(minute)
|
|
if row is None:
|
|
row = buckets[minute] = {}
|
|
cutoff = minute - WINDOW_MINUTES
|
|
for old in [m for m in buckets if m <= cutoff]:
|
|
del buckets[old] # drop minutes that have rolled out of the window
|
|
row[key] = row.get(key, 0) + 1
|
|
|
|
@staticmethod
|
|
def _recent_totals(buckets) -> "dict[str, int]":
|
|
"""Sum each key's count over the last WINDOW minutes (caller holds lock)."""
|
|
cutoff = _now_minute() - WINDOW_MINUTES + 1
|
|
totals: "dict[str, int]" = {}
|
|
for minute, row in buckets.items():
|
|
if minute >= cutoff:
|
|
for key, count in row.items():
|
|
totals[key] = totals.get(key, 0) + count
|
|
return totals
|
|
|
|
def record_inbound(self, category, bucket) -> None:
|
|
with self._lock:
|
|
row = self._inbound.get(category)
|
|
if row is None:
|
|
row = self._inbound[category] = _empty_inbound_row()
|
|
row["total"] += 1
|
|
row[bucket] = row.get(bucket, 0) + 1
|
|
self._bump_recent(self._recent_in, category)
|
|
|
|
def record_outbound(self, source, outcome) -> None:
|
|
with self._lock:
|
|
row = self._outbound.get(source)
|
|
if row is None:
|
|
row = self._outbound[source] = {o: 0 for o in _OUTCOMES}
|
|
row[outcome] = row.get(outcome, 0) + 1
|
|
self._bump_recent(self._recent_out, source)
|
|
|
|
def record_heartbeat(self, name, ts, interval_s) -> None:
|
|
with self._lock:
|
|
self._heartbeats[name] = {"ts": ts, "interval_s": interval_s}
|
|
|
|
def record_event(self, event, referrer, day) -> None:
|
|
with self._lock:
|
|
# Cap distinct referrers so a spray of forged Referer headers can't
|
|
# grow the key space without bound.
|
|
if referrer not in ("direct", "self", REFERRER_OTHER):
|
|
seen = {r for _, r, _ in self._events}
|
|
if referrer not in seen and len(seen) >= MAX_REFERRERS:
|
|
referrer = REFERRER_OTHER
|
|
key = (event, referrer, day)
|
|
self._events[key] = self._events.get(key, 0) + 1
|
|
keep = _recent_days()
|
|
for k in [k for k in self._events if k[2] not in keep]:
|
|
del self._events[k]
|
|
self._bump_recent(self._recent_ev, event)
|
|
|
|
def read(self) -> dict:
|
|
with self._lock:
|
|
inbound = {k: dict(v) for k, v in self._inbound.items()}
|
|
outbound = {k: dict(v) for k, v in self._outbound.items()}
|
|
return {
|
|
"started_at": self.started_at,
|
|
"inbound": inbound,
|
|
"outbound": outbound,
|
|
"inbound_recent": self._recent_totals(self._recent_in),
|
|
"outbound_recent": self._recent_totals(self._recent_out),
|
|
"heartbeats": {k: dict(v) for k, v in self._heartbeats.items()},
|
|
"events": _nest_events(self._events.items()),
|
|
"events_recent": self._recent_totals(self._recent_ev),
|
|
}
|
|
|
|
|
|
# The seven metrics tables, in a stable order. Used by the Postgres init to
|
|
# TRUNCATE (there is no file to `rm`, so we wipe here to keep the since-start
|
|
# semantics) — the guard is exactly this list, and only on the Postgres path.
|
|
_METRICS_TABLES = (
|
|
"meta", "inbound_cume", "outbound_cume", "inbound_minute",
|
|
"outbound_minute", "event_cume", "event_minute",
|
|
)
|
|
|
|
# Postgres advisory-lock key that serializes the once-per-deploy metrics reset
|
|
# (the TRUNCATE below). MUST differ from core/singleton.py's NOTIFIER_LOCK_KEY
|
|
# (1), which the notifier leader holds session-long — they share one advisory
|
|
# lock space, so a collision would deadlock the notifier against the reset.
|
|
_METRICS_BOOT_LOCK_KEY = 0x6D657472 # "metr"
|
|
|
|
|
|
def _boot_token() -> float:
|
|
"""A value identical for every uvicorn worker of one container boot and
|
|
different across deploys, so the since-start reset fires exactly once per
|
|
deployment even though (with the lazy store below) workers connect to
|
|
Postgres seconds or minutes apart. PID 1's start time (`/proc/1/stat` field
|
|
22, ticks since host boot) is shared by all workers in the container and
|
|
changes when the container is recreated. Falls back to a per-process value
|
|
if `/proc` is unavailable (non-Linux) — then each worker resets once, the
|
|
pre-lazy behavior."""
|
|
try:
|
|
with open("/proc/1/stat", "rb") as f:
|
|
# comm (field 2) may contain spaces/parens, so split after the last ')'.
|
|
return float(f.read().rsplit(b")", 1)[1].split()[19])
|
|
except Exception: # noqa: BLE001 - no /proc: degrade to per-process reset
|
|
return float(os.getpid())
|
|
|
|
# Postgres DDL: same shape as the SQLite schema below, but UNLOGGED (fast, non-
|
|
# durable — metrics are since-start and disposable) and one statement per execute
|
|
# (psycopg's extended protocol won't run a multi-statement script).
|
|
_PG_METRICS_SCHEMA = (
|
|
"CREATE UNLOGGED TABLE IF NOT EXISTS meta(key TEXT PRIMARY KEY, value DOUBLE PRECISION)",
|
|
"CREATE UNLOGGED TABLE IF NOT EXISTS inbound_cume("
|
|
"category TEXT, bucket TEXT, count BIGINT, PRIMARY KEY(category, bucket))",
|
|
"CREATE UNLOGGED TABLE IF NOT EXISTS outbound_cume("
|
|
"source TEXT, outcome TEXT, count BIGINT, PRIMARY KEY(source, outcome))",
|
|
"CREATE UNLOGGED TABLE IF NOT EXISTS inbound_minute("
|
|
"minute BIGINT, category TEXT, count BIGINT, PRIMARY KEY(minute, category))",
|
|
"CREATE UNLOGGED TABLE IF NOT EXISTS outbound_minute("
|
|
"minute BIGINT, source TEXT, count BIGINT, PRIMARY KEY(minute, source))",
|
|
"CREATE UNLOGGED TABLE IF NOT EXISTS event_cume("
|
|
"event TEXT, referrer TEXT, day TEXT, count BIGINT, PRIMARY KEY(event, referrer, day))",
|
|
"CREATE UNLOGGED TABLE IF NOT EXISTS event_minute("
|
|
"minute BIGINT, event TEXT, count BIGINT, PRIMARY KEY(minute, event))",
|
|
)
|
|
|
|
|
|
class _SqliteStore:
|
|
"""Cross-process counters in a shared store, so every uvicorn worker tallies into
|
|
one place and the dashboard sees the whole deployment. Two dialects behind one
|
|
interface: a shared SQLite file (WAL) by default, or Postgres (``postgres=True``,
|
|
UNLOGGED tables) when THERMOGRAPH_DATABASE_URL selects it. Either way the tally is
|
|
tiny UPSERTs on a single connection serialized by ``self._lock`` — psycopg
|
|
connections aren't thread-safe, so that lock is what keeps the sharing safe (the
|
|
same role SQLite's ``check_same_thread=False`` plays here).
|
|
"""
|
|
|
|
def __init__(self, path: str, postgres: bool = False) -> None:
|
|
self._lock = threading.Lock()
|
|
self._pg = postgres
|
|
self._ph = "%s" if postgres else "?" # bound-parameter placeholder style
|
|
if postgres:
|
|
import psycopg # local import: only the Postgres path needs it
|
|
# autocommit mirrors SQLite's isolation_level=None: every counter UPSERT
|
|
# lands on its own, and no read leaves an idle-in-transaction.
|
|
self._db = psycopg.connect(path, autocommit=True)
|
|
for stmt in _PG_METRICS_SCHEMA:
|
|
self._db.execute(stmt)
|
|
# SQLite is wiped each boot by an external `ExecStartPre rm` of the file;
|
|
# Postgres has no file to remove, so reset the since-start counters here —
|
|
# but exactly ONCE per deployment, not once per worker. With the lazy store
|
|
# (see _LazyPgStore) workers connect at different times, so an unguarded
|
|
# per-store TRUNCATE would let a late worker wipe counts the others already
|
|
# tallied. Gate it on a per-boot token under an advisory lock: the first
|
|
# worker of a fresh boot finds a stale/absent token and resets; the rest
|
|
# find this boot's token and skip. Guarded to exactly the metrics tables.
|
|
token = _boot_token()
|
|
with self._db.transaction():
|
|
self._db.execute("SELECT pg_advisory_xact_lock(%s)", (_METRICS_BOOT_LOCK_KEY,))
|
|
row = self._db.execute(
|
|
"SELECT value FROM meta WHERE key = 'boot_token'").fetchone()
|
|
if row is None or row[0] != token:
|
|
self._db.execute("TRUNCATE " + ", ".join(_METRICS_TABLES))
|
|
self._db.execute(
|
|
"INSERT INTO meta(key, value) VALUES ('boot_token', %s) "
|
|
"ON CONFLICT (key) DO UPDATE SET value = excluded.value", (token,))
|
|
self._db.execute(
|
|
"INSERT INTO meta(key, value) VALUES ('started_at', %s) "
|
|
"ON CONFLICT (key) DO UPDATE SET value = excluded.value", (time.time(),))
|
|
else:
|
|
# check_same_thread=False: the connection is shared across a worker's threads,
|
|
# serialized by self._lock. isolation_level=None: autocommit each statement.
|
|
self._db = sqlite3.connect(path, check_same_thread=False, isolation_level=None)
|
|
self._db.execute("PRAGMA journal_mode=WAL")
|
|
self._db.execute("PRAGMA synchronous=NORMAL")
|
|
self._db.execute("PRAGMA busy_timeout=3000")
|
|
self._db.executescript(
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS meta(key TEXT PRIMARY KEY, value REAL);
|
|
CREATE TABLE IF NOT EXISTS inbound_cume(
|
|
category TEXT, bucket TEXT, count INTEGER, PRIMARY KEY(category, bucket));
|
|
CREATE TABLE IF NOT EXISTS outbound_cume(
|
|
source TEXT, outcome TEXT, count INTEGER, PRIMARY KEY(source, outcome));
|
|
CREATE TABLE IF NOT EXISTS inbound_minute(
|
|
minute INTEGER, category TEXT, count INTEGER, PRIMARY KEY(minute, category));
|
|
CREATE TABLE IF NOT EXISTS outbound_minute(
|
|
minute INTEGER, source TEXT, count INTEGER, PRIMARY KEY(minute, source));
|
|
CREATE TABLE IF NOT EXISTS event_cume(
|
|
event TEXT, referrer TEXT, day TEXT, count INTEGER,
|
|
PRIMARY KEY(event, referrer, day));
|
|
CREATE TABLE IF NOT EXISTS event_minute(
|
|
minute INTEGER, event TEXT, count INTEGER, PRIMARY KEY(minute, event));
|
|
"""
|
|
)
|
|
# First worker to touch the DB stamps the shared start time; the rest inherit it.
|
|
self._db.execute(
|
|
"INSERT OR IGNORE INTO meta(key, value) VALUES('started_at', ?)", (time.time(),)
|
|
)
|
|
self._last_pruned_minute = -1
|
|
|
|
@property
|
|
def started_at(self) -> float:
|
|
row = self._db.execute(
|
|
"SELECT value FROM meta WHERE key='started_at'").fetchone()
|
|
return row[0] if row else time.time()
|
|
|
|
def _prune(self, minute: int) -> None:
|
|
"""Drop per-minute rows that have rolled out of the window (idempotent across
|
|
workers). Cheap, and only run when the epoch-minute advances."""
|
|
if minute == self._last_pruned_minute:
|
|
return
|
|
self._last_pruned_minute = minute
|
|
cutoff = minute - WINDOW_MINUTES
|
|
self._db.execute(f"DELETE FROM inbound_minute WHERE minute <= {self._ph}", (cutoff,))
|
|
self._db.execute(f"DELETE FROM outbound_minute WHERE minute <= {self._ph}", (cutoff,))
|
|
self._db.execute(f"DELETE FROM event_minute WHERE minute <= {self._ph}", (cutoff,))
|
|
# Event history is kept by day, not by minute.
|
|
placeholders = ",".join([self._ph] * EVENT_DAYS)
|
|
self._db.execute(
|
|
f"DELETE FROM event_cume WHERE day NOT IN ({placeholders})",
|
|
tuple(sorted(_recent_days())),
|
|
)
|
|
|
|
def record_inbound(self, category, bucket) -> None:
|
|
minute = _now_minute()
|
|
with self._lock:
|
|
# Qualify the existing value as <table>.count, not a bare `count`:
|
|
# Postgres rejects the unqualified form in ON CONFLICT DO UPDATE as
|
|
# ambiguous (the target row and the `excluded` pseudo-row both have a
|
|
# `count`), which silently dropped every traffic counter on the
|
|
# Postgres deployments. SQLite accepts the qualified form too.
|
|
self._db.execute(
|
|
f"INSERT INTO inbound_cume(category, bucket, count) VALUES({self._ph}, {self._ph}, 1) "
|
|
"ON CONFLICT(category, bucket) DO UPDATE SET count = inbound_cume.count + 1",
|
|
(category, bucket),
|
|
)
|
|
self._db.execute(
|
|
f"INSERT INTO inbound_minute(minute, category, count) VALUES({self._ph}, {self._ph}, 1) "
|
|
"ON CONFLICT(minute, category) DO UPDATE SET count = inbound_minute.count + 1",
|
|
(minute, category),
|
|
)
|
|
self._prune(minute)
|
|
|
|
def record_outbound(self, source, outcome) -> None:
|
|
minute = _now_minute()
|
|
with self._lock:
|
|
self._db.execute(
|
|
f"INSERT INTO outbound_cume(source, outcome, count) VALUES({self._ph}, {self._ph}, 1) "
|
|
"ON CONFLICT(source, outcome) DO UPDATE SET count = outbound_cume.count + 1",
|
|
(source, outcome),
|
|
)
|
|
self._db.execute(
|
|
f"INSERT INTO outbound_minute(minute, source, count) VALUES({self._ph}, {self._ph}, 1) "
|
|
"ON CONFLICT(minute, source) DO UPDATE SET count = outbound_minute.count + 1",
|
|
(minute, source),
|
|
)
|
|
self._prune(minute)
|
|
|
|
def record_event(self, event, referrer, day) -> None:
|
|
minute = _now_minute()
|
|
with self._lock:
|
|
if referrer not in ("direct", "self", REFERRER_OTHER):
|
|
row = self._db.execute(
|
|
"SELECT COUNT(DISTINCT referrer) FROM event_cume").fetchone()
|
|
known = self._db.execute(
|
|
f"SELECT 1 FROM event_cume WHERE referrer = {self._ph} LIMIT 1", (referrer,)
|
|
).fetchone()
|
|
if not known and row and row[0] >= MAX_REFERRERS:
|
|
referrer = REFERRER_OTHER
|
|
self._db.execute(
|
|
f"INSERT INTO event_cume(event, referrer, day, count) VALUES({self._ph}, {self._ph}, {self._ph}, 1) "
|
|
"ON CONFLICT(event, referrer, day) DO UPDATE SET count = event_cume.count + 1",
|
|
(event, referrer, day),
|
|
)
|
|
self._db.execute(
|
|
f"INSERT INTO event_minute(minute, event, count) VALUES({self._ph}, {self._ph}, 1) "
|
|
"ON CONFLICT(minute, event) DO UPDATE SET count = event_minute.count + 1",
|
|
(minute, event),
|
|
)
|
|
self._prune(minute)
|
|
|
|
def record_heartbeat(self, name, ts, interval_s) -> None:
|
|
# Reuse the meta table (two keys per daemon) so every worker's beats land in the
|
|
# one shared file — the dashboard then sees the notifier's liveness whichever
|
|
# worker answers its poll, even though only the leader worker runs it.
|
|
with self._lock:
|
|
self._db.execute(
|
|
f"INSERT INTO meta(key, value) VALUES({self._ph}, {self._ph}) "
|
|
"ON CONFLICT(key) DO UPDATE SET value = excluded.value",
|
|
(f"hb_ts:{name}", ts),
|
|
)
|
|
self._db.execute(
|
|
f"INSERT INTO meta(key, value) VALUES({self._ph}, {self._ph}) "
|
|
"ON CONFLICT(key) DO UPDATE SET value = excluded.value",
|
|
(f"hb_int:{name}", interval_s),
|
|
)
|
|
|
|
def _read_heartbeats(self) -> "dict[str, dict[str, float]]":
|
|
beats: "dict[str, dict[str, float]]" = {}
|
|
for key, value in self._db.execute(
|
|
"SELECT key, value FROM meta WHERE key LIKE 'hb_ts:%' OR key LIKE 'hb_int:%'"):
|
|
field, _, name = key.partition(":")
|
|
beats.setdefault(name, {})["ts" if field == "hb_ts" else "interval_s"] = value
|
|
return beats
|
|
|
|
def read(self) -> dict:
|
|
cutoff = _now_minute() - WINDOW_MINUTES + 1
|
|
with self._lock:
|
|
inbound: "dict[str, dict[str, int]]" = {}
|
|
for category, bucket, count in self._db.execute(
|
|
"SELECT category, bucket, count FROM inbound_cume"):
|
|
row = inbound.setdefault(category, _empty_inbound_row())
|
|
row[bucket] = row.get(bucket, 0) + count
|
|
row["total"] += count
|
|
# Start from the stable source rows so the dashboard shows them even at zero.
|
|
outbound: "dict[str, dict[str, int]]" = {
|
|
s: {o: 0 for o in _OUTCOMES} for s in SOURCES
|
|
}
|
|
for source, outcome, count in self._db.execute(
|
|
"SELECT source, outcome, count FROM outbound_cume"):
|
|
outbound.setdefault(source, {o: 0 for o in _OUTCOMES})[outcome] = count
|
|
# int() coerces the aggregate: Postgres SUM(BIGINT) returns numeric ->
|
|
# Decimal, which isn't JSON-serializable; SQLite already returns int, so
|
|
# the wrap is a no-op there.
|
|
recent_in = {k: int(v) for k, v in self._db.execute(
|
|
f"SELECT category, SUM(count) FROM inbound_minute WHERE minute >= {self._ph} "
|
|
"GROUP BY category", (cutoff,))}
|
|
recent_out = {k: int(v) for k, v in self._db.execute(
|
|
f"SELECT source, SUM(count) FROM outbound_minute WHERE minute >= {self._ph} "
|
|
"GROUP BY source", (cutoff,))}
|
|
events = _nest_events(
|
|
((event, referrer, day), count)
|
|
for event, referrer, day, count in self._db.execute(
|
|
"SELECT event, referrer, day, count FROM event_cume")
|
|
)
|
|
recent_ev = {k: int(v) for k, v in self._db.execute(
|
|
f"SELECT event, SUM(count) FROM event_minute WHERE minute >= {self._ph} "
|
|
"GROUP BY event", (cutoff,))}
|
|
return {
|
|
"started_at": self.started_at,
|
|
"inbound": inbound,
|
|
"outbound": outbound,
|
|
"inbound_recent": recent_in,
|
|
"outbound_recent": recent_out,
|
|
"heartbeats": self._read_heartbeats(),
|
|
"events": events,
|
|
"events_recent": recent_ev,
|
|
}
|
|
|
|
|
|
def _empty_snapshot() -> dict:
|
|
"""The shape `read()` returns with nothing recorded yet — used while the lazy
|
|
Postgres store hasn't connected, and by `snapshot()` when the store is
|
|
momentarily unreadable, so callers never KeyError on a missing counter."""
|
|
return {"started_at": time.time(), "inbound": {}, "outbound": {},
|
|
"inbound_recent": {}, "outbound_recent": {}, "heartbeats": {},
|
|
"events": {}, "events_recent": {}}
|
|
|
|
|
|
class _LazyPgStore:
|
|
"""Postgres is the intended cross-worker store, but the app and db containers
|
|
start together, so on a cold boot the app imports this module — and used to
|
|
build the store eagerly — before Postgres is accepting connections. That build
|
|
then failed and (see `_make_store`) fell back to a per-process in-memory store
|
|
for the WHOLE container lifetime: metrics fragmented across workers and the
|
|
subscription-notifier heartbeat invisible (it lived in one worker's memory,
|
|
never in the shared table). Fix: don't connect at import. Connect on first use
|
|
— which is always after the app's own DB init, so Postgres is up — and if it
|
|
still isn't ready, retry on a short backoff instead of downgrading forever.
|
|
Once connected, delegate straight through."""
|
|
|
|
_RETRY_INTERVAL_S = 5.0
|
|
|
|
def __init__(self, dsn: str) -> None:
|
|
self._dsn = dsn
|
|
self._real: "_SqliteStore | None" = None
|
|
self._next_try = 0.0
|
|
self._lock = threading.Lock()
|
|
|
|
def _get(self) -> "_SqliteStore | None":
|
|
real = self._real
|
|
if real is not None:
|
|
return real
|
|
now = time.monotonic()
|
|
with self._lock:
|
|
if self._real is not None:
|
|
return self._real
|
|
if now < self._next_try:
|
|
return None
|
|
try:
|
|
self._real = _SqliteStore(self._dsn, postgres=True)
|
|
except Exception: # noqa: BLE001 - Postgres not ready yet; retry later
|
|
self._next_try = now + self._RETRY_INTERVAL_S
|
|
return None
|
|
return self._real
|
|
|
|
# Writes are best-effort until connected — the pre-DB-ready window is seconds,
|
|
# and the module-level record_* already swallow errors. A dropped write there
|
|
# is the same guarantee the old eager store gave once it fell back.
|
|
def record_inbound(self, category, bucket) -> None:
|
|
s = self._get()
|
|
if s is not None:
|
|
s.record_inbound(category, bucket)
|
|
|
|
def record_outbound(self, source, outcome) -> None:
|
|
s = self._get()
|
|
if s is not None:
|
|
s.record_outbound(source, outcome)
|
|
|
|
def record_event(self, event, referrer, day) -> None:
|
|
s = self._get()
|
|
if s is not None:
|
|
s.record_event(event, referrer, day)
|
|
|
|
def record_heartbeat(self, name, ts, interval_s) -> None:
|
|
s = self._get()
|
|
if s is not None:
|
|
s.record_heartbeat(name, ts, interval_s)
|
|
|
|
def read(self) -> dict:
|
|
s = self._get()
|
|
return s.read() if s is not None else _empty_snapshot()
|
|
|
|
|
|
def _make_store():
|
|
"""Pick the store from the environment. A ``postgresql://`` URL in
|
|
THERMOGRAPH_DATABASE_URL routes counters into a shared Postgres store (prod /
|
|
multi-worker) and takes precedence — via `_LazyPgStore`, so a not-yet-ready DB
|
|
at import never permanently downgrades to in-memory. Otherwise a path in
|
|
``THERMOGRAPH_METRICS_DB`` selects the shared SQLite store (needed when running
|
|
multiple workers); otherwise in-memory. Non-Postgres opens still fall back —
|
|
metrics must never break boot."""
|
|
if IS_POSTGRES:
|
|
return _LazyPgStore(_PG_DSN)
|
|
path = os.environ.get("THERMOGRAPH_METRICS_DB", "").strip()
|
|
if path:
|
|
try:
|
|
return _SqliteStore(path)
|
|
except Exception: # noqa: BLE001 - never let instrumentation break startup
|
|
pass
|
|
return _MemStore()
|
|
|
|
|
|
_store = _make_store()
|
|
|
|
|
|
def record_inbound(category: str, status) -> None:
|
|
"""Count one inbound request in ``category``, bucketed by status class."""
|
|
# The metrics endpoint is polled by the dashboard itself — counting it would just
|
|
# show the monitor watching itself, so it's excluded from traffic entirely.
|
|
# The event beacon is likewise not traffic: it is the *record* of traffic, and
|
|
# counting it would double every interaction it reports.
|
|
if category in ("metrics", "event"):
|
|
return
|
|
try:
|
|
_store.record_inbound(category, _status_bucket(status))
|
|
except Exception: # noqa: BLE001 - instrumentation must never break a request
|
|
pass
|
|
|
|
|
|
def record_outbound(phase: str, outcome: str) -> None:
|
|
"""Count one outbound call. ``outcome`` in {ok, retry, error, rate_limited}."""
|
|
try:
|
|
_store.record_outbound(source_for_phase(phase), outcome)
|
|
except Exception: # noqa: BLE001
|
|
pass
|
|
|
|
|
|
# Per-IP token bucket for the unauthenticated event beacon. Bounded: the map is
|
|
# cleared whenever the minute rolls, so it holds at most one minute of clients.
|
|
#
|
|
# 30/min was sized for the four coarse home.* counters (at most a handful per
|
|
# visit). Interaction-level events (chart-metric taps, a comfort slider drag on
|
|
# the compare page) are burstier by an order of magnitude, and a batched beacon
|
|
# spends one token per event in the batch — at 30 the limiter silently truncates
|
|
# ordinary use, which is worse than useless because the loss is invisible and
|
|
# biased toward the most engaged visitors. 120/min is still a hard ceiling on
|
|
# what one IP can write per minute. Env-tunable so it can be tightened on the
|
|
# box without a deploy if the endpoint is ever abused.
|
|
_EVENT_RATE_PER_MIN = int(os.environ.get("THERMOGRAPH_EVENT_RATE_PER_MIN", "120") or 120)
|
|
_rate_lock = threading.Lock()
|
|
_rate_minute = -1
|
|
_rate_counts: "dict[str, int]" = {}
|
|
|
|
|
|
def reset_rate_limiter() -> None:
|
|
"""Drop every token bucket. For tests — one test's beacon traffic must not
|
|
rate-limit the next one's, since TestClient always presents the same IP."""
|
|
global _rate_minute
|
|
with _rate_lock:
|
|
_rate_minute = -1
|
|
_rate_counts.clear()
|
|
|
|
|
|
def rate_ok(ip: "str | None") -> bool:
|
|
"""Public form of the beacon's per-IP token bucket, so a caller that fans one
|
|
event out to several sinks (see web/app.py's ``_record_event``) spends exactly
|
|
one token instead of one per sink. The IP is only ever a key in a map that is
|
|
cleared every minute — it is never persisted anywhere."""
|
|
return _rate_ok(ip)
|
|
|
|
|
|
def _rate_ok(ip: "str | None") -> bool:
|
|
global _rate_minute
|
|
key = ip or "?"
|
|
minute = _now_minute()
|
|
with _rate_lock:
|
|
if minute != _rate_minute:
|
|
_rate_minute = minute
|
|
_rate_counts.clear()
|
|
count = _rate_counts.get(key, 0) + 1
|
|
_rate_counts[key] = count
|
|
return count <= _EVENT_RATE_PER_MIN
|
|
|
|
|
|
def record_event(name: str, referer: "str | None" = None,
|
|
host: "str | None" = None, ip: "str | None" = None,
|
|
rate_limit: bool = True) -> None:
|
|
"""Count one product event. The name is allowlisted and the referrer is taken
|
|
from the request's own Referer header (never client-supplied, which would be
|
|
trivially forgeable). Best-effort and silent: the caller always answers 204,
|
|
so a rejected or rate-limited event tells a prober nothing.
|
|
|
|
``rate_limit=False`` skips the per-IP bucket for a caller that already spent
|
|
the token through ``rate_ok`` — see web/app.py's ``_record_event``."""
|
|
try:
|
|
if rate_limit and not _rate_ok(ip):
|
|
return
|
|
event = normalize_event(name)
|
|
# An unrecognized event carries no referrer dimension — one bucket total.
|
|
referrer = "direct" if event == EVENT_OTHER else normalize_referrer(referer, host)
|
|
_store.record_event(event, referrer, _utc_day())
|
|
except Exception: # noqa: BLE001 - instrumentation must never break a request
|
|
pass
|
|
|
|
|
|
def record_heartbeat(name: str, interval_s: float) -> None:
|
|
"""Stamp ``name``'s liveness. ``interval_s`` is the daemon's expected beat cadence, so
|
|
a reader can tell fresh from stale without knowing the schedule. Best-effort."""
|
|
try:
|
|
_store.record_heartbeat(name, time.time(), interval_s)
|
|
except Exception: # noqa: BLE001 - instrumentation must never break the daemon
|
|
pass
|
|
|
|
|
|
def snapshot() -> dict:
|
|
"""A JSON-serializable copy of all counters plus process facts."""
|
|
try:
|
|
data = _store.read()
|
|
except Exception: # noqa: BLE001 - a locked/broken store must not 500 the dashboard
|
|
data = _empty_snapshot()
|
|
inbound = data["inbound"]
|
|
outbound = data["outbound"]
|
|
started_at = data["started_at"]
|
|
# Convert each heartbeat's absolute timestamp to an age here, on the server clock, so
|
|
# a reader (dashboard) never has to reconcile clock skew to judge freshness.
|
|
now = time.time()
|
|
heartbeats = {
|
|
name: {"age_s": round(now - hb["ts"], 1), "interval_s": hb.get("interval_s")}
|
|
for name, hb in (data.get("heartbeats") or {}).items()
|
|
if hb.get("ts") is not None
|
|
}
|
|
return {
|
|
"pid": os.getpid(),
|
|
"started_at": started_at,
|
|
"uptime_s": round(time.time() - started_at, 1),
|
|
"inbound": inbound,
|
|
"outbound": outbound,
|
|
"inbound_total": sum(v.get("total", 0) for v in inbound.values()),
|
|
"outbound_total": sum(sum(v.values()) for v in outbound.values()),
|
|
# Traffic over the last WINDOW_MINUTES, per category / source.
|
|
"window_minutes": WINDOW_MINUTES,
|
|
"inbound_recent": data["inbound_recent"],
|
|
"outbound_recent": data["outbound_recent"],
|
|
"heartbeats": heartbeats,
|
|
# Product events: {event: {referrer_domain: {day: count}}}, aggregate only.
|
|
"events": data.get("events") or {},
|
|
"events_recent": data.get("events_recent") or {},
|
|
"events_total": sum(
|
|
n
|
|
for refs in (data.get("events") or {}).values()
|
|
for days in refs.values()
|
|
for n in days.values()
|
|
),
|
|
}
|