Files
swp-01-ingest/poller.py
T
irrlichtandClaude Sonnet 5 c43d8496bf Branding auf Wurzelwerk (SWP) umstellen, Deployment systemagnostisch machen
Ersetzt "Lagebild" durch den Projektnamen Wurzelwerk (Stinkwurzpresse) in
allen Units, Configs und Docs. Entfernt konkrete Hostnamen (Bebop, Pi 5,
Debby) aus README/NOTES zugunsten generischer Rollenbezeichnungen.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013vHZCwcCJnz4LKNgvT4bxv
2026-09-04 20:28:57 +02:00

436 lines
16 KiB
Python
Executable File

#!/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, wie bei GDELT
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 15 Minuten abrunden, wie GDELTs Paketgrenzen."""
return ts.replace(second=0, microsecond=0,
minute=(ts.minute // 15) * 15)
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 _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]
aus.append({
"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"])),
})
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, 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 = []
for it in items:
cur = self.db.execute(
"insert or ignore into gesehen (quelle, guid) values (?, ?)",
(quelle["id"], it["guid"]))
if cur.rowcount:
frisch.append(it)
self.db.commit()
if not frisch:
return 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 len(frisch)
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):
neu = 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, roh)
values (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
on conflict (quelle, item_key) do nothing
returning id
""", (quelle["id"], it["guid"], it["url"], it["titel"],
norm_titel(it["titel"]), it["teaser"], it["autor"],
it["pubdate"], slot,
json.dumps(it, ensure_ascii=False, default=str)))
if cur.fetchone():
neu += 1
self.conn.commit()
return neu
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, fehler):
with self.conn.cursor() as cur:
cur.execute("""
insert into poll_laeufe
(quelle, slot, dauer_ms, http_status,
items_gesehen, items_neu, fehler)
values (%s, %s, %s, %s, %s, %s, %s)
""", (quelle, slot, dauer_ms, status, gesehen, neu, 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 = 0
for q in quellen:
if q["typ"] != "rss":
continue
t0 = time.monotonic()
etag, lm = senke.zustand(q["id"])
status = gesehen = neu = 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 = 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, fehler)
ges_neu += neu
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 {dauer:5} ms"
+ (f" {fehler}" if fehler else ""), flush=True)
return ges_gesehen, ges_neu
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 = durchlauf(aktiv, senke, slot)
print(f" = {gesehen} gesehen, {neu} neu", 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())