Der Handler in durchlauf() fing den Fehler einer Quelle ab, damit sie nicht den ganzen Durchlauf killt, liess die abgebrochene Postgres-Transaktion aber stehen. Das notiere() unmittelbar danach lief in genau diese Transaktion und warf InFailedSqlTransaction - ungefangen, also war der Prozess tot. Doppelt aergerlich: der eigentliche Fehler stand nur in `fehler` und wurde weder protokolliert noch ausgegeben, weil beides erst nach dem notiere() kommt. Im Journal blieb allein die Folgewirkung sichtbar. Beide Senken bekommen rollback(), der Fehlerpfad ruft es auf, und notiere() selbst laeuft in einem eigenen try - ein fehlgeschlagenes Protokoll ist ein Grund fuer eine Meldung, nicht fuer einen Abbruch. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013vHZCwcCJnz4LKNgvT4bxv
667 lines
26 KiB
Python
Executable File
667 lines
26 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
|
|
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/",
|
|
"media": "http://search.yahoo.com/mrss/",
|
|
}
|
|
# Header, die wir zum Abruf mitschreiben. Kosten nichts und beantworten
|
|
# spaeter Fragen zur Zwischenspeicherung, die man sonst nicht mehr stellen kann.
|
|
HEADER_MERKEN = ("Date", "Age", "Cache-Control", "Content-Type",
|
|
"Content-Length", "Server", "ETag", "Last-Modified")
|
|
|
|
_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, header) zurueck. Die
|
|
Rohbytes werden bewusst durchgereicht und nicht nur geparst: der Parser
|
|
wird sich aendern, die Historie soll dann rueckwirkend davon profitieren.
|
|
"""
|
|
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)
|
|
header = {k: r.headers.get(k) for k in HEADER_MERKEN
|
|
if r.headers.get(k)}
|
|
return (r.status, roh, r.headers.get("ETag"),
|
|
r.headers.get("Last-Modified"), header)
|
|
except urllib.error.HTTPError as err:
|
|
if err.code == 304:
|
|
return 304, None, etag, last_modified, {}
|
|
raise
|
|
|
|
|
|
_HREF_RE = re.compile(r'href=["\']([^"\']+)["\']', re.I)
|
|
|
|
|
|
def _kategorien(it):
|
|
"""<category>, dc:subject, Atom-term. Vom Verlag selbst vergebene Label."""
|
|
aus = []
|
|
for el in it.findall("category") + it.findall("dc:subject", NS):
|
|
if _text(el):
|
|
aus.append(_text(el))
|
|
for el in it.findall("a:category", NS):
|
|
if el.get("term"):
|
|
aus.append(el.get("term"))
|
|
# Reihenfolge erhalten, Dubletten raus
|
|
return list(dict.fromkeys(aus))
|
|
|
|
|
|
def _medien(it):
|
|
"""enclosure und Media RSS. Acht von zwoelf Quellen liefern Bilder."""
|
|
aus = []
|
|
for el in it.findall("enclosure"):
|
|
if el.get("url"):
|
|
aus.append({"url": el.get("url"), "typ": el.get("type"),
|
|
"bytes": el.get("length"), "rolle": "enclosure"})
|
|
for tag, rolle in (("media:content", "content"),
|
|
("media:thumbnail", "thumbnail")):
|
|
for el in it.findall(tag, NS):
|
|
if el.get("url"):
|
|
aus.append({"url": el.get("url"), "typ": el.get("type"),
|
|
"bytes": el.get("fileSize"), "rolle": rolle})
|
|
return aus
|
|
|
|
|
|
def _links(*html_stuecke):
|
|
"""Verlinkungen aus dem Teaser-HTML, bevor klartext() die Tags entfernt.
|
|
|
|
GDELT haelt das Gegenstueck als PAGE_LINKS im GKG: welche Meldung auf
|
|
welche verweist, ist ein eigenes Signal.
|
|
"""
|
|
aus = []
|
|
for st in html_stuecke:
|
|
if st:
|
|
aus += _HREF_RE.findall(html.unescape(st))
|
|
return list(dict.fromkeys(aus))
|
|
|
|
|
|
def parse_feed(roh):
|
|
"""RSS 2.0, RDF und Atom auf eine gemeinsame Form bringen.
|
|
|
|
Gibt (items, uebersprungen) zurueck. Uebersprungen wird nur, was keinen
|
|
Titel hat - das wird gezaehlt statt stillschweigend verworfen.
|
|
"""
|
|
wurzel = ET.fromstring(roh)
|
|
eintraege = wurzel.findall(".//item") or wurzel.findall(".//a:entry", NS)
|
|
aus, uebersprungen = [], 0
|
|
for pos, it in enumerate(eintraege):
|
|
titel = _erst(it, ["title", "a:title"])
|
|
if not titel:
|
|
uebersprungen += 1
|
|
continue
|
|
|
|
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]
|
|
|
|
# Kurzfassung und Volltext getrennt halten: frueher gewann die kurze
|
|
# description und content:encoded fiel unter den Tisch.
|
|
teaser_html = _erst(it, ["description", "a:summary"])
|
|
voll_html = _erst(it, ["content:encoded", "a:content"])
|
|
if not teaser_html:
|
|
teaser_html = voll_html
|
|
|
|
eintrag = {
|
|
"titel": klartext(titel),
|
|
"url": link,
|
|
"guid": guid,
|
|
"position": pos, # Reihenfolge im Feed = oft die Gewichtung
|
|
"teaser": klartext(teaser_html),
|
|
"volltext": klartext(voll_html) if voll_html is not teaser_html else None,
|
|
"autor": _erst(it, ["author", "dc:creator", "a:author/a:name"]),
|
|
"kategorien": _kategorien(it),
|
|
"medien": _medien(it),
|
|
"links": _links(teaser_html, voll_html),
|
|
"pubdate": datum(_erst(it, ["pubDate", "dc:date", "a:published"])),
|
|
"aktualisiert": datum(_erst(it, ["a:updated", "dc:modified"])),
|
|
}
|
|
if eintrag["pubdate"] is None:
|
|
eintrag["pubdate"] = eintrag["aktualisiert"]
|
|
# Der Hash bleibt bewusst auf Titel/Teaser/URL beschraenkt. Naehme er
|
|
# die neuen Felder auf, gaelte beim ersten Lauf nach dem Update jedes
|
|
# Bestandsitem als revidiert.
|
|
eintrag["inhalt_hash"] = inhalt_hash(
|
|
eintrag["titel"], eintrag["teaser"], eintrag["url"])
|
|
aus.append(eintrag)
|
|
return aus, uebersprungen
|
|
|
|
|
|
# --------------------------------------------------------------------------
|
|
# 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 rollback(self):
|
|
self.db.rollback()
|
|
|
|
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, 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, 0
|
|
|
|
def speichere_feed(self, quelle, slot, status, header, koerper):
|
|
"""Rohes Feed-XML gezippt ablegen, dedupliziert ueber den Hash."""
|
|
if not koerper:
|
|
return False
|
|
h = hashlib.sha256(koerper).hexdigest()
|
|
tag = self.ziel / slot.strftime("%Y%m%d") / "feeds"
|
|
tag.mkdir(parents=True, exist_ok=True)
|
|
pfad = tag / f"{quelle['id']}-{h[:16]}.xml.gz"
|
|
if pfad.exists():
|
|
return False
|
|
pfad.write_bytes(gzip.compress(koerper))
|
|
return True
|
|
|
|
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 rollback(self):
|
|
self.conn.rollback()
|
|
|
|
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 = nachgetragen = 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,
|
|
position, volltext, kategorien, medien, links,
|
|
aktualisiert)
|
|
values (%s, %s, %s, %s, %s, %s, %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,
|
|
position = coalesce(raw_items.position,
|
|
excluded.position),
|
|
volltext = excluded.volltext,
|
|
kategorien = excluded.kategorien,
|
|
medien = excluded.medien,
|
|
links = excluded.links,
|
|
aktualisiert = excluded.aktualisiert,
|
|
inhalt_hash = excluded.inhalt_hash,
|
|
roh = excluded.roh
|
|
where raw_items.inhalt_hash is distinct from excluded.inhalt_hash
|
|
or raw_items.kategorien is null
|
|
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),
|
|
it["position"], it["volltext"], it["kategorien"],
|
|
json.dumps(it["medien"], ensure_ascii=False),
|
|
it["links"], it["aktualisiert"]))
|
|
zeile = cur.fetchone()
|
|
if not zeile:
|
|
continue # unveraendert, nichts zu tun
|
|
rid, ist_neu = zeile
|
|
|
|
# Ob der Inhalt wirklich neu ist, sagt der Unique-Index der
|
|
# Fassungstabelle - zuverlaessiger als ein Vergleich in
|
|
# RETURNING, wo die Zeile schon fortgeschrieben ist.
|
|
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
|
|
returning id
|
|
""", (rid, slot, it["titel"], it["teaser"], it["url"],
|
|
it["inhalt_hash"]))
|
|
inhalt_neu = cur.fetchone() is not None
|
|
|
|
if ist_neu:
|
|
neu += 1
|
|
elif inhalt_neu:
|
|
geaendert += 1
|
|
cur.execute("""
|
|
update raw_items
|
|
set revisionen = revisionen + 1,
|
|
zuletzt_geaendert = now()
|
|
where id = %s
|
|
""", (rid,))
|
|
else:
|
|
nachgetragen += 1 # nur die neuen Felder gefuellt
|
|
self.conn.commit()
|
|
return neu, geaendert, nachgetragen
|
|
|
|
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 speichere_feed(self, quelle, slot, status, header, koerper):
|
|
"""Rohes Feed-XML gezippt ablegen, dedupliziert ueber den Hash.
|
|
|
|
Der Parser wird sich aendern. Nur mit den Rohbytes laesst sich eine
|
|
spaetere Verbesserung rueckwirkend auf die Historie anwenden - dieselbe
|
|
Ueberlegung, aus der raw_items und codings getrennt sind.
|
|
"""
|
|
if not koerper:
|
|
return False
|
|
h = hashlib.sha256(koerper).hexdigest()
|
|
gz = gzip.compress(koerper)
|
|
with self.conn.cursor() as cur:
|
|
cur.execute("""
|
|
insert into feed_abrufe
|
|
(quelle, slot, http_status, header, koerper_hash,
|
|
koerper, bytes_roh, bytes_gz)
|
|
values (%s, %s, %s, %s, %s, %s, %s, %s)
|
|
on conflict (quelle, koerper_hash) do nothing
|
|
returning id
|
|
""", (quelle["id"], slot, status,
|
|
json.dumps(header, ensure_ascii=False), h,
|
|
gz, len(koerper), len(gz)))
|
|
neu = cur.fetchone() is not None
|
|
self.conn.commit()
|
|
return neu
|
|
|
|
def notiere(self, quelle, slot, dauer_ms, status, gesehen, neu,
|
|
geaendert, uebersprungen, bytes_roh, 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, items_uebersprungen,
|
|
bytes_roh, fehler)
|
|
values (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
|
""", (quelle, slot, dauer_ms, status, gesehen, neu, geaendert,
|
|
uebersprungen, bytes_roh, 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 = uebersprungen = bytes_roh = 0
|
|
nachgetragen = 0
|
|
fehler = None
|
|
try:
|
|
status, roh, n_etag, n_lm, header = hole(q["url"], etag, lm)
|
|
if status == 304:
|
|
pass
|
|
else:
|
|
bytes_roh = len(roh)
|
|
senke.speichere_feed(q, slot, status, header, roh)
|
|
items, uebersprungen = parse_feed(roh)
|
|
gesehen = len(items)
|
|
neu, geaendert, nachgetragen = 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)
|
|
# Ohne Rollback bleibt eine abgebrochene Transaktion stehen, und
|
|
# das notiere() gleich darunter scheitert an ihr statt den Fehler
|
|
# festzuhalten - der eigentliche Grund waere dann verloren.
|
|
senke.rollback()
|
|
dauer = int((time.monotonic() - t0) * 1000)
|
|
try:
|
|
senke.notiere(q["id"], slot, dauer, status, gesehen, neu, geaendert,
|
|
uebersprungen, bytes_roh, fehler)
|
|
except Exception as e: # das Protokoll darf den Lauf nicht killen
|
|
senke.rollback()
|
|
print(f" ! {q['id']:14} Protokoll fehlgeschlagen: "
|
|
f"{type(e).__name__}: {e}", flush=True)
|
|
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 "
|
|
f"{uebersprungen:2} uebrg"
|
|
+ (f" {nachgetragen:4} nachgetr" if nachgetragen else "")
|
|
+ f" {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())
|