commit ac7d8c7e6942292216d7a0cf49e99166a82c9dfa Author: irrlicht Date: Fri Sep 4 20:08:22 2026 +0200 Initiales Ingest-Repo (aus gdelt-bot/ingest überführt) Poller, Schema, Quadlet-Units und Quellenregister aus dem Monorepo gdelt-bot in ein eigenes Repo ausgelagert. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01VzNVf1KaHPZ1fzJT2jd18j diff --git a/Containerfile b/Containerfile new file mode 100644 index 0000000..2c7da9d --- /dev/null +++ b/Containerfile @@ -0,0 +1,14 @@ +FROM docker.io/library/python:3.13-slim + +RUN useradd --system --uid 10001 --create-home ingest +WORKDIR /app + +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY poller.py sources.toml ./ +USER ingest + +# Ohne Argumente: Dauerbetrieb im 15-Minuten-Takt, Senke aus $DATABASE_URL. +# -u, damit die Logs unmittelbar im journal landen und nicht im Puffer haengen. +ENTRYPOINT ["python", "-u", "/app/poller.py"] diff --git a/NOTES.md b/NOTES.md new file mode 100644 index 0000000..fd2b6ab --- /dev/null +++ b/NOTES.md @@ -0,0 +1,134 @@ +# Ingest + +Stufe 0 und 1 der Pipeline: Quellenregister und Roh-Erfassung. Keine +Kodierung — was hier landet, soll sich beliebig oft neu kodieren lassen. + +## Dateien + +| Datei | Zweck | +|---|---| +| `sources.toml` | Quellenregister, einzige Wahrheit ueber die Quellen | +| `schema.sql` | Postgres-Schema (`quellen`, `raw_items`, `codings`, `poll_laeufe`) | +| `poller.py` | Poller, zwei Senken: Postgres und gzip-JSONL | +| `Containerfile` | Image fuer Bebop | +| `lagebild-ingest.container`, `lagebild.network` | Quadlet-Units, rootless | +| `ingest.env.example` | Vorlage fuer die Zugangsdaten | + +## Bebop + +```sh +# 1. Schema einspielen +psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql + +# 2. Image bauen +podman build -t localhost/lagebild-ingest:latest . + +# 3. Konfiguration ablegen +mkdir -p ~/.config/lagebild +cp ingest.env.example ~/.config/lagebild/ingest.env +cp sources.toml ~/.config/lagebild/sources.toml +chmod 600 ~/.config/lagebild/ingest.env +$EDITOR ~/.config/lagebild/ingest.env + +# 4. Quadlet installieren +mkdir -p ~/.config/containers/systemd +cp lagebild-ingest.container lagebild.network ~/.config/containers/systemd/ +systemctl --user daemon-reload +systemctl --user start lagebild-ingest +journalctl --user -u lagebild-ingest -f + +# Damit der Dienst ohne offene Sitzung weiterlaeuft: +loginctl enable-linger "$USER" +``` + +`sources.toml` ist eingehaengt, nicht ins Image gebacken: eine neue Quelle ist +ein `systemctl --user restart lagebild-ingest`, kein Neubau. + +## Pi 5 — Schatten-Ingest + +Ohne Postgres, ohne Container: + +```sh +python3 poller.py --sink jsonl:/var/lib/lagebild/schatten +``` + +Legt `/YYYYMMDD/YYYYMMDDHHMMSS.jsonl.gz` an, dedupliziert ueber eine +SQLite daneben. Als `systemd --user`-Unit mit `Restart=always` einrichten. + +Der Sinn ist nicht Redundanz um ihrer selbst willen: RSS-Feeds haben kein +Archiv. Steht Bebops Poller einen Tag still, ist dieser Tag unwiederbringlich +weg. Kodierungen lassen sich wiederholen, Rohdaten nie. + +## Betrieb + +```sh +# Lizenzstatus im Register gegen die aktuelle robots.txt pruefen +python3 poller.py --check-robots # Exit != 0 bei Abweichung + +# ein einzelner Durchlauf +python3 poller.py --sink "$DATABASE_URL" --once +``` + +`--check-robots` gehoert in einen woechentlichen Cron. Ein Verlag, der seine +robots.txt aendert, aendert damit seinen Nutzungsvorbehalt — und das ist der +Punkt, an dem eine Quelle stillgelegt werden muss. + +Aendert sich etwas, wird **`sources.toml` angepasst, nie der Code.** Eine +gesperrte Quelle bleibt im Register stehen (`aktiv = false`), damit +nachvollziehbar bleibt, warum sie fehlt. Der Poller verweigert den Start, +wenn eine als `gesperrt` gefuehrte Quelle aktiv geschaltet ist. + +### Waechter + +```sql +-- aktive Quellen, die seit ueber sechs Stunden nichts Neues liefern +select * from quellen_stillstand; + +-- Fehler der letzten Stunde +select quelle, http_status, fehler, begonnen_am from poll_laeufe +where fehler is not null and begonnen_am > now() - interval '1 hour'; +``` + +### Arbeitsvorrat fuer Yantra + +```sql +select * from unkodiert('qwen3-14b-q4/prompt-1') limit 500; +``` + +Liefert alle Items, die diese Kodierer-Version noch nicht gesehen hat. Eine +neue Prompt-Version ist damit automatisch ein vollstaendiger Recode-Auftrag +ueber die gesamte Historie, ohne dass ein Feed erneut abgerufen wird. + +## Verhalten + +- **15-Minuten-Takt**, an den Slotgrenzen ausgerichtet (+30 s Versatz), + identisch zu GDELTs Paketgrenzen — damit bleiben die Slots vergleichbar. +- **Bedingter GET** ueber ETag/Last-Modified, im Register gespeichert. Im + Test antworteten 4 von 12 Quellen im zweiten Durchlauf mit 304. +- **Dedup** ueber `unique (quelle, item_key)`; `item_key` ist guid, ersatzweise + Link, ersatzweise ein Titel-Hash. +- **Fehlertoleranz**: eine kaputte Quelle beendet den Durchlauf nicht, der + Fehler landet in `poll_laeufe`. +- **Kein Volltext-Abruf.** Nur der Feed selbst wird geholt. Titel und Teaser + reichen als Kodierungsinput, und damit stellt sich die TDM-Frage nach + § 44b UrhG gar nicht erst. + +## Gemessen (2026-09-04) + +Erster Lauf ueber 12 aktive Quellen: 820 Items gesehen, 810 neu, gesamter +Durchlauf unter 3 s. Zweiter Lauf unmittelbar danach: 1 neues Item. + +Das Trigramm-Clustering findet quellenuebergreifende Dubletten sofort: + +``` +zeit | handelsblatt | 1.00 | Bekennerschreiben: Erneut Sabotage am Stromnetz +spiegel | derstandard | 0.92 | Niederlage fuer Donald Trump - Gericht in Missouri +``` + +## Offen + +- Telegram-Fetcher (Telethon/MTProto). Register ist vorbereitet, Eintraege + stehen auf `aktiv = false`. Die Bot-API kann fremde Kanaele nicht lesen. +- Ersatz-URL fuer den BR24-Feed (aktuell 404). +- Quellenpool von 12 auf ~40 erweitern; Regionalzeitungen ueber SearXNG auf + Debby suchen. diff --git a/README.md b/README.md new file mode 100644 index 0000000..13e98bb --- /dev/null +++ b/README.md @@ -0,0 +1,44 @@ +# swp-01-ingest + +Stufe 0/1 der Lagebild-Pipeline: Quellenregister und Roh-Erfassung +(RSS-Poller, Senken Postgres oder gzip-JSONL). Details zu Betrieb, +Verhalten und offenen Punkten stehen in [`NOTES.md`](NOTES.md). + +## Deploy (Bebop, Postgres + Podman/Quadlet) + +```sh +# 1. Schema einspielen +psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql + +# 2. Image bauen +podman build -t localhost/lagebild-ingest:latest . + +# 3. Konfiguration ablegen +mkdir -p ~/.config/lagebild +cp ingest.env.example ~/.config/lagebild/ingest.env +cp sources.toml ~/.config/lagebild/sources.toml +chmod 600 ~/.config/lagebild/ingest.env +$EDITOR ~/.config/lagebild/ingest.env + +# 4. Quadlet installieren +mkdir -p ~/.config/containers/systemd +cp lagebild-ingest.container lagebild.network ~/.config/containers/systemd/ +systemctl --user daemon-reload +systemctl --user start lagebild-ingest +journalctl --user -u lagebild-ingest -f + +# Damit der Dienst ohne offene Sitzung weiterlaeuft: +loginctl enable-linger "$USER" +``` + +`sources.toml` ist eingehaengt, nicht ins Image gebacken: eine neue Quelle +ist ein `systemctl --user restart lagebild-ingest`, kein Neubau. + +## Deploy (Pi 5, ohne Postgres/Container) + +```sh +python3 poller.py --sink jsonl:/var/lib/lagebild/schatten +``` + +Als `systemd --user`-Unit mit `Restart=always` einrichten (siehe +`lagebild-robots.service` / `.timer` fuer den woechentlichen Robots-Check). diff --git a/ingest.env.example b/ingest.env.example new file mode 100644 index 0000000..daeae8b --- /dev/null +++ b/ingest.env.example @@ -0,0 +1,10 @@ +# Nach ~/.config/lagebild/ingest.env kopieren und ausfuellen (chmod 600). +# +# Bei rootless Podman erreichen sich Container ueber ein gemeinsames Netz per +# Containername. Liegt Postgres bereits in einem Netz, dort denselben Namen +# eintragen wie in der Zeile Network= der Quadlet-Datei. +DATABASE_URL=postgresql://lagebild:GEHEIM@lagebild-pg:5432/lagebild + +# Kontaktadresse gehoert in den User-Agent: wer unsere Abrufe sieht, soll +# wissen, wen er anschreiben kann. +INGEST_USER_AGENT=lagebild-ingest/0.1 (+https://DEINE-DOMAIN/lagebild; kontakt@DEINE-DOMAIN) diff --git a/lagebild-ingest.container b/lagebild-ingest.container new file mode 100644 index 0000000..dd8d75a --- /dev/null +++ b/lagebild-ingest.container @@ -0,0 +1,37 @@ +# Quadlet-Unit fuer rootless Podman. +# Ablegen unter ~/.config/containers/systemd/lagebild-ingest.container, +# danach: systemctl --user daemon-reload && systemctl --user start lagebild-ingest +# +# Der Poller laeuft dauerhaft und taktet sich selbst auf die 15-Minuten-Grenzen. +# Deshalb kein systemd-Timer: ein Timer wuerde den bedingten GET-Zustand +# (ETag/Last-Modified) bei jedem Start neu aufbauen muessen. + +[Unit] +Description=Lagebild Ingest-Poller +After=network-online.target +Wants=network-online.target + +[Container] +Image=localhost/lagebild-ingest:latest +ContainerName=lagebild-ingest +AutoUpdate=local + +EnvironmentFile=%h/.config/lagebild/ingest.env + +# Quellenregister von aussen einhaengen: eine neue Quelle ist dann ein +# Neustart, kein Neubau des Images. +Volume=%h/.config/lagebild/sources.toml:/app/sources.toml:ro,Z + +# Dasselbe Netz wie der Postgres-Container. +Network=lagebild.network + +NoNewPrivileges=true +ReadOnly=true +Tmpfs=/tmp + +[Service] +Restart=always +RestartSec=60 + +[Install] +WantedBy=default.target diff --git a/lagebild-robots.service b/lagebild-robots.service new file mode 100644 index 0000000..1e10e34 --- /dev/null +++ b/lagebild-robots.service @@ -0,0 +1,13 @@ +# Wochenpruefung des Lizenzstatus. Nach ~/.config/systemd/user/ ablegen. +# Kein Quadlet: der Lauf ist kurzlebig und soll bei Abweichung fehlschlagen, +# damit systemd den Fehlerzustand sichtbar haelt. + +[Unit] +Description=Lagebild - robots.txt der Quellen pruefen + +[Service] +Type=oneshot +ExecStart=/usr/bin/podman run --rm \ + --network=host \ + -v %h/.config/lagebild/sources.toml:/app/sources.toml:ro \ + localhost/lagebild-ingest:latest --check-robots diff --git a/lagebild-robots.timer b/lagebild-robots.timer new file mode 100644 index 0000000..41a9077 --- /dev/null +++ b/lagebild-robots.timer @@ -0,0 +1,13 @@ +# Nach ~/.config/systemd/user/ ablegen, dann: +# systemctl --user enable --now lagebild-robots.timer + +[Unit] +Description=Lagebild - woechentliche robots.txt-Pruefung + +[Timer] +OnCalendar=Mon 05:30 +RandomizedDelaySec=30m +Persistent=true + +[Install] +WantedBy=timers.target diff --git a/lagebild.network b/lagebild.network new file mode 100644 index 0000000..26916bb --- /dev/null +++ b/lagebild.network @@ -0,0 +1,6 @@ +# Ablegen unter ~/.config/containers/systemd/lagebild.network +[Unit] +Description=Lagebild internes Netz + +[Network] +NetworkName=lagebild diff --git a/poller.py b/poller.py new file mode 100755 index 0000000..1014b45 --- /dev/null +++ b/poller.py @@ -0,0 +1,435 @@ +#!/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://... Bebop, die eigentliche Datenbank + --sink jsonl:/pfad/dir Pi 5, 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", + "lagebild-ingest/0.1 (+https://example.org/lagebild; kontakt@example.org)", +) +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 Bebop 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()) diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..1d24c5b --- /dev/null +++ b/requirements.txt @@ -0,0 +1 @@ +psycopg[binary]>=3.2 diff --git a/schema.sql b/schema.sql new file mode 100644 index 0000000..e883db0 --- /dev/null +++ b/schema.sql @@ -0,0 +1,126 @@ +-- Schema fuer den Ingest. +-- +-- Grundregel: raw_items ist append-only und wird nie ueberschrieben. +-- Kodierungen sind eine getrennte Stufe mit eigener Versionierung, damit die +-- gesamte Historie neu kodiert werden kann, ohne einen Feed erneut abzurufen. + +create extension if not exists pg_trgm; + +-- -------------------------------------------------------------------------- +-- Quellen: aus ingest/sources.toml synchronisiert, nie von Hand pflegen. +-- -------------------------------------------------------------------------- +create table if not exists quellen ( + id text primary key, + name text not null, + typ text not null check (typ in ('rss', 'telegram')), + url text not null, + land text not null, + traeger text, + ressort text, + lizenz_status text not null + check (lizenz_status in ('offen','gesperrt','ungeprueft')), + robots_geprueft_am date, + aktiv boolean not null default false, + notiz text, + + -- Poller-Zustand, nicht aus der TOML + etag text, + last_modified text, + zuletzt_geprueft timestamptz, + zuletzt_geaendert timestamptz +); + +comment on column quellen.zuletzt_geprueft is 'letzter Abruf, auch wenn 304'; +comment on column quellen.zuletzt_geaendert is 'letzter Abruf mit neuen Items'; + +-- -------------------------------------------------------------------------- +-- Rohdaten. Append-only. RSS-Feeds haben kein Archiv: was hier fehlt, ist +-- unwiederbringlich weg. +-- -------------------------------------------------------------------------- +create table if not exists raw_items ( + id bigserial primary key, + quelle text not null references quellen(id), + item_key text not null, -- guid/link, nur innerhalb der Quelle eindeutig + url text, + titel text not null, + titel_norm text not null, -- normalisiert, fuer Clustering + teaser text, + autor text, + pubdate timestamptz, -- Angabe der Quelle, kann fehlen/luegen + abgerufen_am timestamptz not null default now(), + slot timestamptz not null, -- auf 15 min abgerundet, GDELT-kompatibel + sprache text not null default 'de', + roh jsonb not null, -- alle geparsten Feldwerte, unveraendert + + unique (quelle, item_key) +); + +create index if not exists raw_items_slot_idx on raw_items (slot desc); +create index if not exists raw_items_pubdate_idx on raw_items (pubdate desc nulls last); +create index if not exists raw_items_quelle_idx on raw_items (quelle, abgerufen_am desc); +-- Titel-Clustering laeuft ueber Trigramm-Aehnlichkeit in SQL, nicht in Python +create index if not exists raw_items_titelnorm_trgm + on raw_items using gin (titel_norm gin_trgm_ops); + +-- -------------------------------------------------------------------------- +-- Kodierungen. Eine Zeile pro (Item, Kodierer-Version). Die Nutzlast bleibt +-- absichtlich jsonb, solange nicht gemessen ist, bis zu welcher CAMEO-Ebene +-- das gewaehlte Modell traegt. +-- -------------------------------------------------------------------------- +create table if not exists codings ( + id bigserial primary key, + raw_item_id bigint not null references raw_items(id) on delete cascade, + coder_version text not null, -- z.B. 'qwen3-14b-q4/prompt-3' + kodiert_am timestamptz not null default now(), + ebene text not null -- erreichte CAMEO-Ebene + check (ebene in ('keine','quad','root','base','event')), + konfidenz real, + nutzlast jsonb not null, + + unique (raw_item_id, coder_version) +); + +create index if not exists codings_version_idx on codings (coder_version, kodiert_am desc); + +-- Arbeitsvorrat fuer Yantra: was diese Kodierer-Version noch nicht gesehen hat. +-- Aufruf: select * from unkodiert('qwen3-14b-q4/prompt-3') limit 500; +create or replace function unkodiert(v text) +returns setof raw_items language sql stable as $$ + select r.* from raw_items r + where not exists ( + select 1 from codings c + where c.raw_item_id = r.id and c.coder_version = v + ) + order by r.id +$$; + +-- -------------------------------------------------------------------------- +-- Betriebsprotokoll. Damit der Watchdog eine stillstehende Quelle bemerkt, +-- bevor Tage fehlen. +-- -------------------------------------------------------------------------- +create table if not exists poll_laeufe ( + id bigserial primary key, + quelle text not null references quellen(id), + slot timestamptz not null, + begonnen_am timestamptz not null default now(), + dauer_ms integer, + http_status integer, -- 304 = unveraendert + items_gesehen integer not null default 0, + items_neu integer not null default 0, + fehler text +); + +create index if not exists poll_laeufe_slot_idx on poll_laeufe (slot desc); +create index if not exists poll_laeufe_quelle_idx on poll_laeufe (quelle, begonnen_am desc); + +-- Welche aktive Quelle hat seit ueber sechs Stunden nichts Neues geliefert? +create or replace view quellen_stillstand as +select q.id, q.name, max(r.abgerufen_am) as letztes_item, + now() - max(r.abgerufen_am) as stille +from quellen q +left join raw_items r on r.quelle = q.id +where q.aktiv +group by q.id, q.name +having max(r.abgerufen_am) is null + or now() - max(r.abgerufen_am) > interval '6 hours' +order by stille desc nulls first; diff --git a/sources.toml b/sources.toml new file mode 100644 index 0000000..2883cad --- /dev/null +++ b/sources.toml @@ -0,0 +1,218 @@ +# Quellenregister +# +# Einzige Wahrheit ueber die Quellen. Der Poller synchronisiert die Tabelle +# `quellen` bei jedem Start aus dieser Datei. +# +# lizenz_status: +# offen - robots.txt erlaubt den Abruf, Feed wird gepollt +# gesperrt - robots.txt oder AGB verbieten den Abruf, Feed wird NICHT gepollt +# ungeprueft - noch nicht geprueft, wird NICHT gepollt +# +# `aktiv = false` schaltet eine Quelle ab, ohne sie aus dem Register zu loeschen. +# Geloescht wird nie: die Historie soll nachvollziehbar bleiben. +# +# Geprueft am 2026-09-04 mit User-Agent aus poller.py (siehe USER_AGENT). +# Neupruefung: `python poller.py --check-robots` + +[[quelle]] +id = "tagesschau" +name = "tagesschau.de" +typ = "rss" +url = "https://www.tagesschau.de/index~rss2.xml" +land = "DE" +traeger = "oeffentlich-rechtlich" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "spiegel" +name = "Der Spiegel" +typ = "rss" +url = "https://www.spiegel.de/schlagzeilen/index.rss" +land = "DE" +traeger = "privat" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "zeit" +name = "Die Zeit" +typ = "rss" +url = "https://newsfeed.zeit.de/index" +land = "DE" +traeger = "privat" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "faz" +name = "Frankfurter Allgemeine Zeitung" +typ = "rss" +url = "https://www.faz.net/rss/aktuell/" +land = "DE" +traeger = "privat" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "sz" +name = "Sueddeutsche Zeitung" +typ = "rss" +url = "https://rss.sueddeutsche.de/rss/Topthemen" +land = "DE" +traeger = "privat" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "welt" +name = "Die Welt" +typ = "rss" +url = "https://www.welt.de/feeds/latest.rss" +land = "DE" +traeger = "privat" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "taz" +name = "die tageszeitung" +typ = "rss" +url = "https://taz.de/!p4608;rss/" +land = "DE" +traeger = "genossenschaft" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "ntv" +name = "n-tv" +typ = "rss" +url = "https://www.n-tv.de/rss" +land = "DE" +traeger = "privat" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "heise" +name = "heise online" +typ = "rss" +url = "https://www.heise.de/rss/heise-atom.xml" +land = "DE" +traeger = "privat" +ressort = "technik" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "dlf" +name = "Deutschlandfunk" +typ = "rss" +url = "https://www.deutschlandfunk.de/nachrichten-100.rss" +land = "DE" +traeger = "oeffentlich-rechtlich" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "handelsblatt" +name = "Handelsblatt" +typ = "rss" +url = "https://www.handelsblatt.com/contentexport/feed/schlagzeilen" +land = "DE" +traeger = "privat" +ressort = "wirtschaft" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +[[quelle]] +id = "derstandard" +name = "Der Standard" +typ = "rss" +url = "https://www.derstandard.at/rss" +land = "AT" +traeger = "privat" +ressort = "allgemein" +lizenz_status = "offen" +robots_geprueft_am = 2026-09-04 +aktiv = true + +# --- gesperrt: robots.txt verbietet den Abruf des Feed-Pfades --------------- + +[[quelle]] +id = "fr" +name = "Frankfurter Rundschau" +typ = "rss" +url = "https://www.fr.de/rssfeed.rdf" +land = "DE" +traeger = "privat" +ressort = "allgemein" +lizenz_status = "gesperrt" +robots_geprueft_am = 2026-09-04 +aktiv = false +notiz = "robots.txt: Disallow fuer den Feed-Pfad (geprueft 2026-09-04)" + +[[quelle]] +id = "nzz" +name = "Neue Zuercher Zeitung" +typ = "rss" +url = "https://www.nzz.ch/recent.rss" +land = "CH" +traeger = "privat" +ressort = "allgemein" +lizenz_status = "gesperrt" +robots_geprueft_am = 2026-09-04 +aktiv = false +notiz = "robots.txt: Disallow fuer den Feed-Pfad (geprueft 2026-09-04)" + +# --- kaputt / offen --------------------------------------------------------- + +[[quelle]] +id = "br24" +name = "BR24" +typ = "rss" +url = "https://www.br.de/nachrichten/rss/alle-meldungen.xml" +land = "DE" +traeger = "oeffentlich-rechtlich" +ressort = "allgemein" +lizenz_status = "ungeprueft" +aktiv = false +notiz = "URL liefert 404 (2026-09-04), Ersatz-URL suchen" + +# --- Telegram: Register vorbereitet, Fetcher noch nicht implementiert ------- +# Die Bot-API kann fremde Kanaele nicht lesen. Noetig ist MTProto (Telethon) +# mit eigener api_id/api_hash von my.telegram.org. `t.me/s/` waere die +# Notloesung, ist aber HTML-Scraping und bricht bei Markup-Aenderungen. + +[[quelle]] +id = "tg_tagesschau" +name = "tagesschau (Telegram)" +typ = "telegram" +url = "tagesschau" +land = "DE" +traeger = "oeffentlich-rechtlich" +ressort = "allgemein" +lizenz_status = "ungeprueft" +aktiv = false +notiz = "Fetcher fehlt noch (Telethon)"