"""The directory's SQLite store. One row per site, keyed by registrable domain, plus the submission and check history the rate limits and the recheck schedule read. Client addresses are never stored: a submission keeps a salted hash, and the salt is replaced every month so old hashes stop being comparable to new ones. """ from datetime import UTC, datetime, timedelta import hashlib import os import secrets import sqlite3 SCHEMA_VERSION = "1" RECHECK_DAYS = 7 RETRY_HOURS = 24 TRANSIENT_LIMIT = 3 RETENTION_DAYS = 90 # Rate limits. Listing is slower than checking because it writes to the # directory; checking only costs a fetch. IP_LISTINGS_PER_HOUR = 5 IP_LISTINGS_PER_DAY = 20 IP_CHECKS_PER_HOUR = 20 DOMAIN_COOLDOWN_MINUTES = 10 DOMAIN_REJECTS_BEFORE_SLOWDOWN = 3 DOMAIN_SLOW_COOLDOWN_HOURS = 24 LISTINGS_PER_HOUR = 60 SCHEMA = """ CREATE TABLE IF NOT EXISTS meta ( key TEXT PRIMARY KEY, value TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS sites ( id INTEGER PRIMARY KEY, domain TEXT NOT NULL UNIQUE, url TEXT NOT NULL, title TEXT NOT NULL, description TEXT, language TEXT, state TEXT NOT NULL CHECK (state IN ('listed', 'removed')), reason TEXT, listed_at TEXT NOT NULL, last_check_at TEXT, last_ok_at TEXT, next_check_at TEXT, transient_failures INTEGER NOT NULL DEFAULT 0, etag TEXT, last_modified TEXT ); CREATE INDEX IF NOT EXISTS sites_due ON sites(next_check_at) WHERE state = 'listed'; CREATE INDEX IF NOT EXISTS sites_state ON sites(state, domain); CREATE TABLE IF NOT EXISTS submissions ( id INTEGER PRIMARY KEY, created_at TEXT NOT NULL, url TEXT NOT NULL, domain TEXT, ip_hash TEXT NOT NULL, outcome TEXT NOT NULL CHECK (outcome IN ( 'listed', 'updated', 'rejected', 'unreachable', 'rate_limited', 'blocked', 'checked')), site_id INTEGER REFERENCES sites(id) ON DELETE SET NULL ); CREATE INDEX IF NOT EXISTS submissions_ip ON submissions(ip_hash, created_at); CREATE INDEX IF NOT EXISTS submissions_domain ON submissions(domain, created_at); CREATE TABLE IF NOT EXISTS checks ( id INTEGER PRIMARY KEY, site_id INTEGER REFERENCES sites(id) ON DELETE CASCADE, checked_at TEXT NOT NULL, kind TEXT NOT NULL CHECK (kind IN ('submission', 'recheck')), result TEXT NOT NULL CHECK (result IN ('pass', 'fail', 'unreachable')), http_status INTEGER, bytes INTEGER, duration_ms INTEGER ); CREATE INDEX IF NOT EXISTS checks_site ON checks(site_id, checked_at DESC); CREATE TABLE IF NOT EXISTS findings ( check_id INTEGER NOT NULL REFERENCES checks(id) ON DELETE CASCADE, section TEXT NOT NULL, level TEXT NOT NULL CHECK (level IN ('must', 'should')), code TEXT NOT NULL, message TEXT NOT NULL, location TEXT ); CREATE INDEX IF NOT EXISTS findings_check ON findings(check_id); CREATE TABLE IF NOT EXISTS blocked ( domain TEXT PRIMARY KEY, blocked_at TEXT NOT NULL, reason TEXT NOT NULL ); """ def now() -> datetime: """Return the current time, in UTC.""" return datetime.now(UTC) def path() -> str: """Return where the database lives.""" return os.environ.get("MEWS_DB", "mews.db") def connect(database: str | None = None) -> sqlite3.Connection: """Open the database with the settings every caller wants.""" connection = sqlite3.connect(database or path(), timeout=5.0) connection.row_factory = sqlite3.Row connection.execute("PRAGMA journal_mode = WAL") connection.execute("PRAGMA foreign_keys = ON") connection.execute("PRAGMA busy_timeout = 5000") return connection def init(connection: sqlite3.Connection) -> None: """Create the schema if it isn't there yet.""" connection.executescript(SCHEMA) connection.execute( "INSERT OR IGNORE INTO meta (key, value) VALUES ('schema_version', ?)", (SCHEMA_VERSION,), ) connection.commit() def ip_hash(connection: sqlite3.Connection, address: str) -> str: """Hash a client address with a salt that is replaced each month.""" stamp = now().strftime("%Y-%m") row = connection.execute( "SELECT value FROM meta WHERE key = 'ip_salt_month'" ).fetchone() if row is None or row["value"] != stamp: connection.execute( "INSERT INTO meta (key, value) VALUES ('ip_salt', ?) " "ON CONFLICT(key) DO UPDATE SET value = excluded.value", (secrets.token_hex(16),), ) connection.execute( "INSERT INTO meta (key, value) VALUES ('ip_salt_month', ?) " "ON CONFLICT(key) DO UPDATE SET value = excluded.value", (stamp,), ) connection.commit() salt = connection.execute( "SELECT value FROM meta WHERE key = 'ip_salt'" ).fetchone()["value"] return hashlib.sha256(f"{salt}{address}".encode()).hexdigest() def _since(hours: float) -> str: return (now() - timedelta(hours=hours)).isoformat() def _count(connection: sqlite3.Connection, sql: str, *args: object) -> int: return connection.execute(sql, args).fetchone()[0] def blocked_reason(connection: sqlite3.Connection, domain: str) -> str | None: """Return why a domain was removed for good, or None if it wasn't.""" row = connection.execute( "SELECT reason FROM blocked WHERE domain = ?", (domain,) ).fetchone() return row["reason"] if row else None def rate_limited( connection: sqlite3.Connection, *, client: str, domain: str, listing: bool ) -> str | None: """Return why this request can't go ahead now, as a sentence for the author. Checking a page is cheap and gets a looser limit than listing a site. """ if not listing: checks = _count( connection, "SELECT COUNT(*) FROM submissions WHERE ip_hash = ? AND created_at > ?", client, _since(1), ) if checks >= IP_CHECKS_PER_HOUR: return "You've checked several pages already. Try again in an hour." return None if ( _count( connection, "SELECT COUNT(*) FROM submissions WHERE ip_hash = ? AND created_at > ? " "AND outcome IN ('listed', 'updated', 'rejected')", client, _since(1), ) >= IP_LISTINGS_PER_HOUR ): return "You've submitted several sites already. Try again in an hour." if ( _count( connection, "SELECT COUNT(*) FROM submissions WHERE ip_hash = ? AND created_at > ? " "AND outcome IN ('listed', 'updated', 'rejected')", client, _since(24), ) >= IP_LISTINGS_PER_DAY ): return "You've submitted several sites already. Try again tomorrow." rejects = _count( connection, "SELECT COUNT(*) FROM submissions WHERE domain = ? AND outcome = 'rejected' " "AND created_at > ?", domain, _since(24), ) cooldown = ( DOMAIN_SLOW_COOLDOWN_HOURS if rejects >= DOMAIN_REJECTS_BEFORE_SLOWDOWN else DOMAIN_COOLDOWN_MINUTES / 60 ) # Only a real attempt at listing counts. Checking a page is the sensible # thing to do first, and it must not lock the author out of listing it. if _count( connection, "SELECT COUNT(*) FROM submissions WHERE domain = ? AND created_at > ? " "AND outcome IN ('listed', 'updated', 'rejected')", domain, _since(cooldown), ): if cooldown > 1: return ( f"{domain} didn't pass the checks a few times today. Fix the " "page and try again tomorrow." ) return f"{domain} was submitted a moment ago. Try again in ten minutes." if ( _count( connection, "SELECT COUNT(*) FROM submissions WHERE created_at > ? " "AND outcome IN ('listed', 'updated')", _since(1), ) >= LISTINGS_PER_HOUR ): return "The directory is busy right now. Try again in an hour." return None def record_submission( connection: sqlite3.Connection, *, url: str, domain: str | None, client: str, outcome: str, site_id: int | None = None, ) -> int: """Store one submission and return its id.""" cursor = connection.execute( "INSERT INTO submissions (created_at, url, domain, ip_hash, outcome, site_id) " "VALUES (?, ?, ?, ?, ?, ?)", (now().isoformat(), url, domain, client, outcome, site_id), ) connection.commit() return int(cursor.lastrowid or 0) def record_check( connection: sqlite3.Connection, *, site_id: int | None, kind: str, result: str, findings: list[tuple[str, str, str, str, str]] = [], # noqa: B006 http_status: int | None = None, size: int | None = None, duration_ms: int | None = None, ) -> int: """Store one check and its findings, returning the check id.""" cursor = connection.execute( "INSERT INTO checks (site_id, checked_at, kind, result, http_status, bytes, " "duration_ms) VALUES (?, ?, ?, ?, ?, ?, ?)", (site_id, now().isoformat(), kind, result, http_status, size, duration_ms), ) check_id = int(cursor.lastrowid or 0) connection.executemany( "INSERT INTO findings (check_id, section, level, code, message, location) " "VALUES (?, ?, ?, ?, ?, ?)", [(check_id, *finding) for finding in findings], ) connection.commit() return check_id def next_check(moment: datetime | None = None) -> str: """Return when a site that just passed should be looked at again. The interval is jittered so submissions made together do not come back as a burst a week later. """ base = (moment or now()) + timedelta(days=RECHECK_DAYS) offset = secrets.randbelow(24 * 60) - 12 * 60 return (base + timedelta(minutes=offset)).isoformat() def upsert_site( connection: sqlite3.Connection, *, domain: str, url: str, title: str, description: str, language: str, etag: str | None, last_modified: str | None, ) -> tuple[int, bool]: """List a site, or refresh the row of one already listed. Returns the site id and whether the row already existed. """ moment = now().isoformat() row = connection.execute( "SELECT id FROM sites WHERE domain = ?", (domain,) ).fetchone() if row is None: cursor = connection.execute( "INSERT INTO sites (domain, url, title, description, language, state, " "listed_at, last_check_at, last_ok_at, next_check_at, etag, last_modified) " "VALUES (?, ?, ?, ?, ?, 'listed', ?, ?, ?, ?, ?, ?)", ( domain, url, title, description, language, moment, moment, moment, next_check(), etag, last_modified, ), ) connection.commit() return int(cursor.lastrowid or 0), False connection.execute( "UPDATE sites SET url = ?, title = ?, description = ?, language = ?, " "state = 'listed', reason = NULL, last_check_at = ?, last_ok_at = ?, " "next_check_at = ?, transient_failures = 0, etag = ?, last_modified = ? " "WHERE id = ?", ( url, title, description, language, moment, moment, next_check(), etag, last_modified, row["id"], ), ) connection.commit() return int(row["id"]), True def listed(connection: sqlite3.Connection) -> list[sqlite3.Row]: """Return every listed site, in the order the directory page shows them.""" return list( connection.execute( "SELECT * FROM sites WHERE state = 'listed' ORDER BY domain" ).fetchall() ) def due(connection: sqlite3.Connection, limit: int = 50) -> list[sqlite3.Row]: """Return the listed sites whose next check has come round.""" return list( connection.execute( "SELECT * FROM sites WHERE state = 'listed' AND next_check_at <= ? " "ORDER BY next_check_at LIMIT ?", (now().isoformat(), limit), ).fetchall() ) def site(connection: sqlite3.Connection, domain: str) -> sqlite3.Row | None: """Return one site by registrable domain.""" return connection.execute( "SELECT * FROM sites WHERE domain = ?", (domain,) ).fetchone() def mark_pass(connection: sqlite3.Connection, site_id: int, **columns: object) -> None: """Record a check a site passed.""" moment = now().isoformat() connection.execute( "UPDATE sites SET last_check_at = ?, last_ok_at = ?, next_check_at = ?, " "transient_failures = 0, etag = ?, last_modified = ? WHERE id = ?", ( moment, moment, next_check(), columns.get("etag"), columns.get("last_modified"), site_id, ), ) connection.commit() def mark_transient(connection: sqlite3.Connection, site_id: int) -> int: """Count one failure that wasn't the page's fault, and return the new count. A site is only dropped once these pile up: a weekend of downtime should not empty the directory. """ connection.execute( "UPDATE sites SET transient_failures = transient_failures + 1, " "last_check_at = ?, next_check_at = ? WHERE id = ?", ( now().isoformat(), (now() + timedelta(hours=RETRY_HOURS)).isoformat(), site_id, ), ) connection.commit() return int( connection.execute( "SELECT transient_failures FROM sites WHERE id = ?", (site_id,) ).fetchone()[0] ) def remove(connection: sqlite3.Connection, site_id: int, reason: str) -> None: """Drop a site from the directory, keeping the row for the record.""" connection.execute( "UPDATE sites SET state = 'removed', reason = ?, last_check_at = ?, " "next_check_at = NULL WHERE id = ?", (reason, now().isoformat(), site_id), ) connection.commit() def block(connection: sqlite3.Connection, domain: str, reason: str) -> None: """Refuse future submissions from a domain.""" connection.execute( "INSERT INTO blocked (domain, blocked_at, reason) VALUES (?, ?, ?) " "ON CONFLICT(domain) DO UPDATE SET reason = excluded.reason", (domain, now().isoformat(), reason), ) connection.commit() def unblock(connection: sqlite3.Connection, domain: str) -> None: """Allow submissions from a domain again.""" connection.execute("DELETE FROM blocked WHERE domain = ?", (domain,)) connection.commit() def last_findings(connection: sqlite3.Connection, site_id: int) -> list[sqlite3.Row]: """Return the findings from a site's most recent check.""" row = connection.execute( "SELECT id FROM checks WHERE site_id = ? ORDER BY checked_at DESC LIMIT 1", (site_id,), ).fetchone() if row is None: return [] return list( connection.execute( "SELECT * FROM findings WHERE check_id = ?", (row["id"],) ).fetchall() ) def prune(connection: sqlite3.Connection, days: int = RETENTION_DAYS) -> None: """Delete submission and check history past the retention window.""" cutoff = (now() - timedelta(days=days)).isoformat() connection.execute("DELETE FROM submissions WHERE created_at < ?", (cutoff,)) connection.execute("DELETE FROM checks WHERE checked_at < ?", (cutoff,)) connection.commit()