#!/usr/bin/env python3 """Ingest-Poller fuer den deutschen Ereignisdatensatz. Liest die Quellen aus sources.toml, holt alle 15 Minuten die Feeds und legt jedes neue Item roh ab. Bewusst ohne Kodierung: was hier landet, soll sich beliebig oft neu kodieren lassen. Zwei Senken: --sink postgres://... Produktivbetrieb, die eigentliche Datenbank --sink jsonl:/pfad/dir Schatten-Ingest als gzip-JSONL pro Slot Aufrufe: poller.py --sink "$DATABASE_URL" # Dauerbetrieb, 15-Minuten-Takt poller.py --sink "$DATABASE_URL" --once # ein Durchlauf poller.py --check-robots # Lizenzstatus nachpruefen """ import argparse import datetime as dt import gzip import hashlib import html import json import os import pathlib import re import signal import sys import time import tomllib import sqlite3 import urllib.error import urllib.parse import urllib.request import urllib.robotparser import xml.etree.ElementTree as ET from email.utils import parsedate_to_datetime HIER = pathlib.Path(__file__).parent USER_AGENT = os.environ.get( "INGEST_USER_AGENT", "wurzelwerk-ingest/0.1 (+https://stinkwurzpresse.de/wurzelwerk; kontakt@stinkwurzpresse.de)", ) TAKT = 15 * 60 # Slotlaenge in Sekunden VERSATZ = 30 # Sekunden nach der Slotgrenze, nicht exakt darauf NS = { "a": "http://www.w3.org/2005/Atom", "dc": "http://purl.org/dc/elements/1.1/", "content": "http://purl.org/rss/1.0/modules/content/", } _lauf = True def _stop(signum, frame): global _lauf _lauf = False print(f"[{_jetzt():%H:%M:%S}] Signal {signum}, beende nach diesem Durchlauf", flush=True) # -------------------------------------------------------------------------- # Hilfsfunktionen # -------------------------------------------------------------------------- def _jetzt(): return dt.datetime.now(dt.timezone.utc) def slot_von(ts): """Auf TAKT-Sekunden seit Epoch abrunden (bei 15 min auf :00/:15/:30/:45).""" sekunden = int(ts.timestamp()) return dt.datetime.fromtimestamp( (sekunden // TAKT) * TAKT, tz=dt.timezone.utc) def norm_titel(t): """Wie build_site.norm_title: Basis fuer das Trigramm-Clustering.""" t = (t or "").lower().split(" - ")[0].split(" | ")[0] # \w haelt Umlaute und ss, ohne die es im Deutschen nicht geht return re.sub(r"\W+", " ", t, flags=re.UNICODE).strip() _TAG_RE = re.compile(r"<[^>]+>") def klartext(s): """HTML aus Teasern entfernen. Feeds kodieren teils doppelt.""" if not s: return None s = html.unescape(html.unescape(s)) s = _TAG_RE.sub(" ", s) return re.sub(r"\s+", " ", s).strip() or None def datum(s): s = (s or "").strip() if not s: return None try: d = parsedate_to_datetime(s) except Exception: try: d = dt.datetime.fromisoformat(s.replace("Z", "+00:00")) except Exception: return None if d.tzinfo is None: d = d.replace(tzinfo=dt.timezone.utc) return d.astimezone(dt.timezone.utc) def inhalt_hash(titel, teaser, url): """Fingerabdruck des Item-Inhalts. Identisch zur SQL-Formel in der Migration, damit Bestandsdaten und neue Zeilen denselben Hash tragen. md5 reicht - es geht um Gleichheit, nicht um Faelschungssicherheit. """ roh = "\x1f".join((titel or "", teaser or "", url or "")) return hashlib.md5(roh.encode("utf-8"), usedforsecurity=False).hexdigest() def _text(el): return (el.text or "").strip() if el is not None and el.text else None def _erst(item, pfade): for p in pfade: el = item.find(p, NS) if ":" in p else item.find(p) v = _text(el) if v: return v return None # -------------------------------------------------------------------------- # Feed-Abruf und -Parsing # -------------------------------------------------------------------------- def hole(url, etag=None, last_modified=None, timeout=30): """Bedingter GET. Gibt (status, bytes|None, etag, last_modified) zurueck.""" kopf = {"User-Agent": USER_AGENT, "Accept-Encoding": "gzip"} if etag: kopf["If-None-Match"] = etag if last_modified: kopf["If-Modified-Since"] = last_modified req = urllib.request.Request(url, headers=kopf) try: with urllib.request.urlopen(req, timeout=timeout) as r: roh = r.read() if r.headers.get("Content-Encoding") == "gzip": roh = gzip.decompress(roh) return (r.status, roh, r.headers.get("ETag"), r.headers.get("Last-Modified")) except urllib.error.HTTPError as err: if err.code == 304: return 304, None, etag, last_modified raise def parse_feed(roh): """RSS 2.0, RDF und Atom auf eine gemeinsame Form bringen.""" wurzel = ET.fromstring(roh) eintraege = wurzel.findall(".//item") or wurzel.findall(".//a:entry", NS) aus = [] for it in eintraege: titel = _erst(it, ["title", "a:title"]) if not titel: continue # ohne Titel ist ein Item fuer uns wertlos link = _erst(it, ["link", "a:id"]) if not link: for le in it.findall("a:link", NS): if le.get("rel") in (None, "alternate") and le.get("href"): link = le.get("href") break guid = _erst(it, ["guid", "a:id"]) or link if not guid: guid = hashlib.sha256(titel.encode()).hexdigest()[:32] eintrag = { "titel": klartext(titel), "url": link, "guid": guid, "teaser": klartext(_erst(it, ["description", "content:encoded", "a:summary", "a:content"])), "autor": _erst(it, ["author", "dc:creator", "a:author/a:name"]), "pubdate": datum(_erst(it, ["pubDate", "dc:date", "a:published", "a:updated"])), } eintrag["inhalt_hash"] = inhalt_hash( eintrag["titel"], eintrag["teaser"], eintrag["url"]) aus.append(eintrag) return aus # -------------------------------------------------------------------------- # Senken # -------------------------------------------------------------------------- class JsonlSenke: """Schatten-Ingest fuer den Pi: gzip-JSONL, ein Verzeichnis pro Tag. Kein Schema, keine Abfragen - nur die Garantie, dass die Rohdaten noch da sind, wenn der Produktiv-Poller einen Tag lang stillstand. Die SQLite daneben haelt nur die gesehenen Schluessel, damit nicht bei jedem Slot der komplette Feed erneut geschrieben wird. """ def __init__(self, ziel): self.ziel = pathlib.Path(ziel) self.ziel.mkdir(parents=True, exist_ok=True) self.db = sqlite3.connect(self.ziel / "gesehen.sqlite") self.db.execute("create table if not exists gesehen (" "quelle text, guid text, hash text, " "primary key (quelle, guid))") self.db.commit() def sync_quellen(self, quellen): (self.ziel / "sources.snapshot.json").write_text( json.dumps(quellen, indent=2, ensure_ascii=False, default=str)) def zustand(self, qid): return None, None def schreibe(self, quelle, items, slot): frisch, neu, geaendert = [], 0, 0 for it in items: row = self.db.execute( "select hash from gesehen where quelle = ? and guid = ?", (quelle["id"], it["guid"])).fetchone() if row is None: self.db.execute("insert into gesehen (quelle, guid, hash) " "values (?, ?, ?)", (quelle["id"], it["guid"], it["inhalt_hash"])) neu += 1 elif row[0] != it["inhalt_hash"]: self.db.execute("update gesehen set hash = ? " "where quelle = ? and guid = ?", (it["inhalt_hash"], quelle["id"], it["guid"])) geaendert += 1 else: continue frisch.append(it) self.db.commit() if not frisch: return 0, 0 tag = self.ziel / slot.strftime("%Y%m%d") tag.mkdir(exist_ok=True) pfad = tag / (slot.strftime("%Y%m%d%H%M%S") + ".jsonl.gz") # gzip-Member lassen sich anhaengen, das Ergebnis bleibt gueltiges gzip with gzip.open(pfad, "at", encoding="utf-8") as f: for it in frisch: f.write(json.dumps( {"quelle": quelle["id"], "slot": slot.isoformat(), "abgerufen_am": _jetzt().isoformat(), **it}, ensure_ascii=False, default=str) + "\n") return neu, geaendert def notiere(self, *a, **kw): pass def close(self): self.db.close() class PostgresSenke: def __init__(self, dsn): import psycopg self.conn = psycopg.connect(dsn, autocommit=False) def sync_quellen(self, quellen): with self.conn.cursor() as cur: for q in quellen: cur.execute(""" insert into quellen (id, name, typ, url, land, traeger, ressort, lizenz_status, robots_geprueft_am, aktiv, notiz) values (%(id)s, %(name)s, %(typ)s, %(url)s, %(land)s, %(traeger)s, %(ressort)s, %(lizenz_status)s, %(robots_geprueft_am)s, %(aktiv)s, %(notiz)s) on conflict (id) do update set name = excluded.name, typ = excluded.typ, url = excluded.url, land = excluded.land, traeger = excluded.traeger, ressort = excluded.ressort, lizenz_status = excluded.lizenz_status, robots_geprueft_am = excluded.robots_geprueft_am, aktiv = excluded.aktiv, notiz = excluded.notiz """, {k: q.get(k) for k in ("id", "name", "typ", "url", "land", "traeger", "ressort", "lizenz_status", "robots_geprueft_am", "aktiv", "notiz")}) self.conn.commit() def zustand(self, qid): with self.conn.cursor() as cur: cur.execute("select etag, last_modified from quellen where id = %s", (qid,)) r = cur.fetchone() return (r[0], r[1]) if r else (None, None) def schreibe(self, quelle, items, slot): """Neue Items einfuegen, geaenderte fortschreiben. Redaktionen aendern Ueberschriften nach der Veroeffentlichung. Frueher lief das in ein `do nothing` und war unsichtbar. Jetzt schreibt der Upsert den aktuellen Stand fort und legt jede Fassung zusaetzlich in raw_item_versionen ab - auch die erste, damit die Historie vollstaendig ist. Bleibt der Hash gleich, liefert das WHERE keine Zeile zurueck und es passiert nichts. """ neu = geaendert = 0 with self.conn.cursor() as cur: for it in items: cur.execute(""" insert into raw_items (quelle, item_key, url, titel, titel_norm, teaser, autor, pubdate, slot, inhalt_hash, roh) values (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) on conflict (quelle, item_key) do update set url = excluded.url, titel = excluded.titel, titel_norm = excluded.titel_norm, teaser = excluded.teaser, autor = excluded.autor, pubdate = excluded.pubdate, inhalt_hash = excluded.inhalt_hash, roh = excluded.roh, revisionen = raw_items.revisionen + 1, zuletzt_geaendert = now() where raw_items.inhalt_hash is distinct from excluded.inhalt_hash returning id, (xmax = 0) as ist_neu """, (quelle["id"], it["guid"], it["url"], it["titel"], norm_titel(it["titel"]), it["teaser"], it["autor"], it["pubdate"], slot, it["inhalt_hash"], json.dumps(it, ensure_ascii=False, default=str))) zeile = cur.fetchone() if not zeile: continue # unveraendert rid, ist_neu = zeile if ist_neu: neu += 1 else: geaendert += 1 cur.execute(""" insert into raw_item_versionen (raw_item_id, slot, titel, teaser, url, inhalt_hash) values (%s, %s, %s, %s, %s, %s) on conflict (raw_item_id, inhalt_hash) do nothing """, (rid, slot, it["titel"], it["teaser"], it["url"], it["inhalt_hash"])) self.conn.commit() return neu, geaendert def merke_zustand(self, qid, etag, last_modified, geaendert): with self.conn.cursor() as cur: cur.execute(""" update quellen set etag = %s, last_modified = %s, zuletzt_geprueft = now(), zuletzt_geaendert = case when %s then now() else zuletzt_geaendert end where id = %s """, (etag, last_modified, geaendert, qid)) self.conn.commit() def notiere(self, quelle, slot, dauer_ms, status, gesehen, neu, geaendert, fehler): with self.conn.cursor() as cur: cur.execute(""" insert into poll_laeufe (quelle, slot, dauer_ms, http_status, items_gesehen, items_neu, items_geaendert, fehler) values (%s, %s, %s, %s, %s, %s, %s, %s) """, (quelle, slot, dauer_ms, status, gesehen, neu, geaendert, fehler)) self.conn.commit() def close(self): self.conn.close() # -------------------------------------------------------------------------- # Durchlauf # -------------------------------------------------------------------------- def lade_quellen(pfad): with open(pfad, "rb") as f: return tomllib.load(f)["quelle"] def durchlauf(quellen, senke, slot): ges_neu = ges_gesehen = ges_geaendert = 0 for q in quellen: if q["typ"] != "rss": continue t0 = time.monotonic() etag, lm = senke.zustand(q["id"]) status = gesehen = neu = geaendert = 0 fehler = None try: status, roh, n_etag, n_lm = hole(q["url"], etag, lm) if status == 304: pass else: items = parse_feed(roh) gesehen = len(items) neu, geaendert = senke.schreibe(q, items, slot) if hasattr(senke, "merke_zustand"): senke.merke_zustand(q["id"], n_etag, n_lm, neu > 0) except Exception as e: # eine Quelle darf nicht fehler = f"{type(e).__name__}: {e}" # den Durchlauf killen status = getattr(e, "code", None) dauer = int((time.monotonic() - t0) * 1000) senke.notiere(q["id"], slot, dauer, status, gesehen, neu, geaendert, fehler) ges_neu += neu ges_geaendert += geaendert ges_gesehen += gesehen marke = "!" if fehler else (" " if status != 304 else ".") print(f" {marke} {q['id']:14} {str(status):>4} " f"{gesehen:4} gesehen {neu:4} neu {geaendert:3} rev {dauer:5} ms" + (f" {fehler}" if fehler else ""), flush=True) return ges_gesehen, ges_neu, ges_geaendert def check_robots(quellen): """Lizenzstatus gegen die aktuelle robots.txt nachpruefen.""" print(f"{'quelle':14} {'im Register':12} {'robots.txt':12} Abweichung") abweichungen = 0 for q in quellen: if q["typ"] != "rss": continue p = urllib.parse.urlparse(q["url"]) rp = urllib.robotparser.RobotFileParser() rp.set_url(f"{p.scheme}://{p.netloc}/robots.txt") try: rp.read() ist = "offen" if rp.can_fetch(USER_AGENT, q["url"]) else "gesperrt" except Exception as e: ist = "fehler" soll = q["lizenz_status"] # 'ungeprueft' im Register ist keine Abweichung, sondern eine Aufgabe weicht = soll != "ungeprueft" and ist != soll abweichungen += weicht print(f"{q['id']:14} {soll:12} {ist:12} {'<<< PRUEFEN' if weicht else ''}") print(f"\n{abweichungen} Abweichung(en). " f"Bei Treffern sources.toml anpassen, nicht den Code.") return abweichungen def main(): ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) ap.add_argument("--sources", default=str(HIER / "sources.toml")) ap.add_argument("--sink", default=os.environ.get("DATABASE_URL"), help="postgres://... oder jsonl:/pfad/zum/verzeichnis") ap.add_argument("--once", action="store_true", help="ein Durchlauf, dann Ende") ap.add_argument("--check-robots", action="store_true", help="Lizenzstatus nachpruefen, nichts abrufen") args = ap.parse_args() alle = lade_quellen(args.sources) if args.check_robots: return 1 if check_robots(alle) else 0 if not args.sink: ap.error("--sink oder DATABASE_URL noetig") aktiv = [q for q in alle if q.get("aktiv")] gesperrt = [q for q in alle if q.get("lizenz_status") == "gesperrt"] if any(q.get("aktiv") for q in gesperrt): sys.exit("FEHLER: gesperrte Quelle ist aktiv geschaltet") if args.sink.startswith("jsonl:"): senke = JsonlSenke(args.sink[len("jsonl:"):]) else: senke = PostgresSenke(args.sink) senke.sync_quellen(alle) print(f"{len(aktiv)} aktive Quellen, {len(gesperrt)} gesperrt, " f"Senke: {args.sink.split('@')[-1]}", flush=True) signal.signal(signal.SIGTERM, _stop) signal.signal(signal.SIGINT, _stop) try: while _lauf: slot = slot_von(_jetzt()) print(f"[{slot:%Y-%m-%d %H:%M}] Slot", flush=True) gesehen, neu, geaendert = durchlauf(aktiv, senke, slot) print(f" = {gesehen} gesehen, {neu} neu, " f"{geaendert} revidiert", flush=True) if args.once: break ziel = slot + dt.timedelta(seconds=TAKT + VERSATZ) while _lauf and _jetzt() < ziel: time.sleep(min(5, (ziel - _jetzt()).total_seconds())) finally: senke.close() return 0 if __name__ == "__main__": sys.exit(main())