""" Snapshot and baseline management. Snapshots are timestamped crawl results. The baseline is a manually approved reference state. Compromised content must never flow into the baseline automatically — every update requires an explicit `approve` call. """ import csv import hashlib import json import logging import sqlite3 from datetime import datetime, timezone from pathlib import Path from typing import Optional logger = logging.getLogger(__name__) def url_key(url: str) -> str: """Stable filename-safe key for a URL (first 16 hex chars of SHA-256).""" return hashlib.sha256(url.encode("utf-8")).hexdigest()[:16] class BaselineManager: def __init__(self, data_dir: str | Path): self.data_dir = Path(data_dir) self.baseline_dir = self.data_dir / "baseline" self.pages_dir = self.baseline_dir / "pages" self.snapshots_dir = self.data_dir / "snapshots" # ------------------------------------------------------------------ # Snapshot operations # ------------------------------------------------------------------ def save_snapshot(self, pages: list[dict], errors: list[dict] | None = None) -> Path: """ Persist a crawl+extraction result as a timestamped snapshot. Crawl errors (network failures) are recorded in the manifest so they survive into the comparison/report stage. Returns the snapshot directory path. """ ts = _utc_stamp() snap_dir = self.snapshots_dir / ts snap_pages_dir = snap_dir / "pages" snap_pages_dir.mkdir(parents=True, exist_ok=True) manifest = { "timestamp": ts, "page_count": len(pages), "urls": [p["url"] for p in pages], "errors": errors or [], } (snap_dir / "manifest.json").write_text( json.dumps(manifest, indent=2, ensure_ascii=False), encoding="utf-8" ) for page in pages: fname = url_key(page["url"]) + ".json" (snap_pages_dir / fname).write_text( json.dumps(page, indent=2, ensure_ascii=False), encoding="utf-8" ) logger.info("Snapshot saved: %s (%d pages)", snap_dir, len(pages)) return snap_dir def latest_snapshot_dir(self) -> Optional[Path]: """Return the most recent snapshot directory, or None.""" if not self.snapshots_dir.exists(): return None dirs = sorted( (d for d in self.snapshots_dir.iterdir() if d.is_dir()), key=lambda d: d.name, ) return dirs[-1] if dirs else None def load_snapshot(self, snap_dir: Optional[Path] = None) -> dict: """ Load a snapshot (default: latest). Returns {manifest, pages: {url: page_dict}, errors: list, dir: Path}. """ if snap_dir is None: snap_dir = self.latest_snapshot_dir() if snap_dir is None or not snap_dir.exists(): return {} manifest = json.loads((snap_dir / "manifest.json").read_text(encoding="utf-8")) pages: dict[str, dict] = {} pages_dir = snap_dir / "pages" for f in pages_dir.glob("*.json"): page = json.loads(f.read_text(encoding="utf-8")) pages[page["url"]] = page return { "manifest": manifest, "pages": pages, "errors": manifest.get("errors", []), "dir": snap_dir, } # ------------------------------------------------------------------ # Baseline operations # ------------------------------------------------------------------ def baseline_exists(self) -> bool: return (self.baseline_dir / "manifest.json").exists() def load_baseline(self) -> dict: """ Load the current approved baseline. Returns {manifest, pages: {url: page_dict}}. """ if not self.baseline_exists(): return {"manifest": {}, "pages": {}} manifest = json.loads( (self.baseline_dir / "manifest.json").read_text(encoding="utf-8") ) pages: dict[str, dict] = {} if self.pages_dir.exists(): for f in self.pages_dir.glob("*.json"): page = json.loads(f.read_text(encoding="utf-8")) pages[page["url"]] = page return {"manifest": manifest, "pages": pages} # ------------------------------------------------------------------ # Approve operations (the only way baseline content is ever updated) # ------------------------------------------------------------------ def approve_all( self, snap_dir: Optional[Path] = None, note: str = "", approved_by: str = "cli", ) -> list[str]: """ Approve every page in the snapshot as the new baseline. Existing baseline pages NOT present in snapshot are kept unless --rebuild is requested (caller responsibility to wipe baseline_dir first). """ snap = self.load_snapshot(snap_dir) if not snap: raise RuntimeError("No snapshot found to approve.") self.pages_dir.mkdir(parents=True, exist_ok=True) approved: list[str] = [] for url, page in snap["pages"].items(): fname = url_key(url) + ".json" (self.pages_dir / fname).write_text( json.dumps(page, indent=2, ensure_ascii=False), encoding="utf-8" ) approved.append(url) self._write_manifest(approved, note, approved_by, snap.get("dir")) logger.info("Approved %d URLs into baseline.", len(approved)) return approved def approve_url( self, url: str, snap_dir: Optional[Path] = None, note: str = "", approved_by: str = "cli", ) -> bool: """Approve a single URL from snapshot into the baseline.""" snap = self.load_snapshot(snap_dir) if url not in snap.get("pages", {}): logger.error("URL not found in snapshot: %s", url) return False self.pages_dir.mkdir(parents=True, exist_ok=True) fname = url_key(url) + ".json" (self.pages_dir / fname).write_text( json.dumps(snap["pages"][url], indent=2, ensure_ascii=False), encoding="utf-8", ) # Update manifest incrementally if self.baseline_exists(): manifest = json.loads( (self.baseline_dir / "manifest.json").read_text(encoding="utf-8") ) else: manifest = {"urls": []} if url not in manifest.get("urls", []): manifest.setdefault("urls", []).append(url) manifest["last_updated"] = datetime.now(timezone.utc).isoformat() manifest["last_note"] = note manifest["last_approved_by"] = approved_by (self.baseline_dir / "manifest.json").write_text( json.dumps(manifest, indent=2, ensure_ascii=False), encoding="utf-8" ) logger.info("Approved single URL into baseline: %s", url) return True def rebuild_baseline( self, snap_dir: Optional[Path] = None, note: str = "", approved_by: str = "cli", ) -> list[str]: """ Wipe the entire baseline and replace it with the given snapshot. Use when a major legitimate redesign has taken place. """ import shutil if self.baseline_dir.exists(): shutil.rmtree(self.baseline_dir) self.pages_dir.mkdir(parents=True, exist_ok=True) return self.approve_all(snap_dir=snap_dir, note=note, approved_by=approved_by) # ------------------------------------------------------------------ # Asset hash baseline # ------------------------------------------------------------------ @property def _asset_hashes_path(self) -> Path: return self.baseline_dir / "asset_hashes.json" def load_asset_hashes(self) -> dict: """Load approved asset hashes. Returns {} if none approved yet.""" if not self._asset_hashes_path.exists(): return {} return json.loads(self._asset_hashes_path.read_text(encoding="utf-8")) def save_asset_hashes(self, hashes: dict) -> None: """Persist asset hashes as the new approved state.""" self.baseline_dir.mkdir(parents=True, exist_ok=True) self._asset_hashes_path.write_text( json.dumps(hashes, indent=2, ensure_ascii=False), encoding="utf-8" ) logger.info("Asset hashes saved: %d entries", len(hashes)) # ------------------------------------------------------------------ # AI content inventory (SQLite store: verdict + content description + url→hash) # ------------------------------------------------------------------ # # EINE Quelle der Wahrheit pro Site: jedes geprüfte Objekt steht mit Inhalts-Hash, # URL(s), Sicherheits-Verdikt UND neutraler KI-Inhaltsangabe in der DB. Der Analyzer # arbeitet weiterhin mit dem In-Memory-Dict {"entries": {hash:…}, "url_hashes": {url:hash}}; # nur die Persistenz liegt in SQLite (abfragbar, skaliert, CSV-Export möglich). @property def _content_db_path(self) -> Path: return self.data_dir / "content_inventory.db" def _connect_inventory(self) -> sqlite3.Connection: self.data_dir.mkdir(parents=True, exist_ok=True) conn = sqlite3.connect(self._content_db_path) conn.row_factory = sqlite3.Row conn.execute(""" CREATE TABLE IF NOT EXISTS objects ( hash TEXT PRIMARY KEY, kind TEXT, url TEXT, description TEXT, category TEXT, severity TEXT, confidence REAL, explanation TEXT, model TEXT, dismissed INTEGER DEFAULT 0, checked_at TEXT, first_seen TEXT, last_seen TEXT )""") conn.execute(""" CREATE TABLE IF NOT EXISTS object_urls ( url TEXT PRIMARY KEY, hash TEXT NOT NULL, last_seen TEXT )""") conn.execute("CREATE INDEX IF NOT EXISTS idx_object_urls_hash ON object_urls(hash)") return conn def load_ai_ledger(self) -> dict: """Load the AI store into the in-memory dict the analyzer uses. Returns {"entries": {hash: {...}}, "url_hashes": {url: hash}}.""" self._migrate_json_ledger() conn = self._connect_inventory() try: entries: dict = {} for row in conn.execute("SELECT * FROM objects"): d = dict(row) h = d.pop("hash") d.pop("first_seen", None) d.pop("last_seen", None) d["dismissed"] = bool(d.get("dismissed")) entries[h] = d url_hashes = {r["url"]: r["hash"] for r in conn.execute("SELECT url, hash FROM object_urls")} finally: conn.close() return {"entries": entries, "url_hashes": url_hashes} def save_ai_ledger(self, ledger: dict) -> None: """Persist the in-memory dict back to SQLite. Objekte werden geupsertet (first_seen bleibt erhalten); object_urls wird ersetzt (damit Entfernungen wirken).""" now = datetime.now(timezone.utc).isoformat() conn = self._connect_inventory() try: for h, e in ledger.get("entries", {}).items(): conn.execute(""" INSERT INTO objects (hash, kind, url, description, category, severity, confidence, explanation, model, dismissed, checked_at, first_seen, last_seen) VALUES (:hash,:kind,:url,:description,:category,:severity,:confidence, :explanation,:model,:dismissed,:checked_at,:now,:now) ON CONFLICT(hash) DO UPDATE SET kind=excluded.kind, url=excluded.url, description=excluded.description, category=excluded.category, severity=excluded.severity, confidence=excluded.confidence, explanation=excluded.explanation, model=excluded.model, dismissed=excluded.dismissed, checked_at=excluded.checked_at, last_seen=excluded.last_seen """, { "hash": h, "kind": e.get("kind"), "url": e.get("url"), "description": e.get("description", ""), "category": e.get("category"), "severity": e.get("severity"), "confidence": e.get("confidence"), "explanation": e.get("explanation", ""), "model": e.get("model"), "dismissed": 1 if e.get("dismissed") else 0, "checked_at": e.get("checked_at"), "now": now, }) conn.execute("DELETE FROM object_urls") conn.executemany( "INSERT INTO object_urls (url, hash, last_seen) VALUES (?,?,?)", [(u, h, now) for u, h in ledger.get("url_hashes", {}).items()]) conn.commit() finally: conn.close() def _migrate_json_ledger(self) -> None: """Einmalige Migration eines alten ai_ledger.json nach SQLite.""" old = self.data_dir / "ai_ledger.json" if not old.exists(): return try: data = json.loads(old.read_text(encoding="utf-8")) except (json.JSONDecodeError, ValueError): old.rename(old.with_suffix(".json.corrupt")) return self.save_ai_ledger({"entries": data.get("entries", {}), "url_hashes": data.get("url_hashes", {})}) old.rename(old.with_suffix(".json.migrated")) logger.info("ai_ledger.json nach SQLite migriert (%d Objekte).", len(data.get("entries", {}))) def query_inventory(self, search: str | None = None, kind: str | None = None) -> list[dict]: """Inventar abfragen (für das `inventory`-Kommando).""" self._migrate_json_ledger() conn = self._connect_inventory() try: sql, params, cond = "SELECT * FROM objects", [], [] if kind: cond.append("kind = ?"); params.append(kind) if search: cond.append("(description LIKE ? OR url LIKE ? OR category LIKE ?)") params += [f"%{search}%"] * 3 if cond: sql += " WHERE " + " AND ".join(cond) sql += " ORDER BY kind, url" return [dict(r) for r in conn.execute(sql, params)] finally: conn.close() def export_inventory_csv(self, path: str | Path) -> int: """Inventar als CSV exportieren. Returns Anzahl Zeilen.""" cols = ["hash", "kind", "url", "description", "category", "severity", "confidence", "explanation", "model", "dismissed", "checked_at", "first_seen", "last_seen"] rows = self.query_inventory() with Path(path).open("w", newline="", encoding="utf-8") as f: w = csv.DictWriter(f, fieldnames=cols, extrasaction="ignore") w.writeheader() for r in rows: w.writerow(r) return len(rows) # ------------------------------------------------------------------ # AI model-router state (circuit-breaker telemetry + /models catalog) # ------------------------------------------------------------------ @property def _ai_router_state_path(self) -> Path: return self.data_dir / "ai_router_state.json" def load_ai_router_state(self) -> dict: """Load the model-router state. Returns a safe default if none/corrupt.""" default = {"models": {}, "catalog": {"slugs": [], "fetched_at": None}} if not self._ai_router_state_path.exists(): return default try: data = json.loads(self._ai_router_state_path.read_text(encoding="utf-8")) except (json.JSONDecodeError, ValueError): return default data.setdefault("models", {}) data.setdefault("catalog", {"slugs": [], "fetched_at": None}) return data def save_ai_router_state(self, state: dict) -> None: """Persist the model-router state.""" self.data_dir.mkdir(parents=True, exist_ok=True) self._ai_router_state_path.write_text( json.dumps(state, indent=2, ensure_ascii=False), encoding="utf-8" ) def invalidate_ai_url_hashes(self, urls: list[str]) -> int: """Entfernt URLs aus der url_hashes-Erinnerung. Folge: Die KI lädt diese Bilder beim nächsten Lauf erneut und bewertet sie inhaltlich neu (geänderte Bytes → neuer Hash → Cache-Miss → Analyse). Wird von der Datei-Prüfung (check-assets) für geänderte Bild-Assets aufgerufen. Returns Anzahl entfernter Einträge.""" if not urls: return 0 ledger = self.load_ai_ledger() uh = ledger.get("url_hashes", {}) removed = [u for u in urls if u in uh] for u in removed: del uh[u] if removed: self.save_ai_ledger(ledger) return len(removed) def dismiss_ai_entries(self, fingerprints: list[str] | None = None) -> int: """ Mark ledger entries as dismissed (false-positive acknowledgement). None = dismiss all currently flagged entries. Returns count dismissed. """ ledger = self.load_ai_ledger() entries = ledger.get("entries", {}) count = 0 for fp, entry in entries.items(): if fingerprints is not None and fp not in fingerprints: continue if entry.get("category", "clean") != "clean" and not entry.get("dismissed"): entry["dismissed"] = True count += 1 if count: self.save_ai_ledger(ledger) return count # ------------------------------------------------------------------ # Run-state (drives auto-scheduling of the weekly checks) # ------------------------------------------------------------------ @property def _state_path(self) -> Path: return self.data_dir / "state.json" def load_state(self) -> dict: """Load the run-state file. Returns {'last_runs': {}} if none exists.""" if not self._state_path.exists(): return {"last_runs": {}} try: data = json.loads(self._state_path.read_text(encoding="utf-8")) except (json.JSONDecodeError, ValueError): return {"last_runs": {}} data.setdefault("last_runs", {}) return data def save_state(self, state: dict) -> None: """Persist the run-state file.""" self.data_dir.mkdir(parents=True, exist_ok=True) self._state_path.write_text( json.dumps(state, indent=2, ensure_ascii=False), encoding="utf-8" ) def is_check_due(self, name: str, interval_days: int) -> bool: """True if a periodic check has never run or last ran >= interval_days ago.""" last = self.load_state().get("last_runs", {}).get(name) if not last: return True try: last_dt = datetime.fromisoformat(last) except ValueError: return True age = datetime.now(timezone.utc) - last_dt return age.total_seconds() >= interval_days * 86400 def mark_check_run(self, name: str) -> None: """Record that a periodic check just ran (now, UTC).""" state = self.load_state() state["last_runs"][name] = datetime.now(timezone.utc).isoformat() self.save_state(state) # ------------------------------------------------------------------ # Internal helpers # ------------------------------------------------------------------ def _write_manifest( self, urls: list[str], note: str, approved_by: str, source_snap: Optional[Path], ) -> None: self.baseline_dir.mkdir(parents=True, exist_ok=True) manifest = { "approved_at": datetime.now(timezone.utc).isoformat(), "approved_by": approved_by, "note": note, "source_snapshot": str(source_snap) if source_snap else "", "urls": urls, } (self.baseline_dir / "manifest.json").write_text( json.dumps(manifest, indent=2, ensure_ascii=False), encoding="utf-8" ) def _utc_stamp() -> str: return datetime.now(timezone.utc).strftime("%Y%m%d_%H%M%S")