Der Abruf warf bisher fast alles weg, was im XML stand. Neu ausgewertet werden Kategorien (<category>, dc:subject, Atom-term), content:encoded zusaetzlich zur kurzen description, Medien-URLs aus enclosure und Media RSS, Verweise aus dem Teaser-HTML und die Position im Feed beim ersten Sehen. Die Kategorien sind fuer Stufe 2 der wertvollste Posten: vom Verlag selbst vergebene Themenlabel, brauchbar als Vorfilter und als Prompt-Kontext. Gemessen an zwoelf Quellen tragen 189 von 814 Items Kategorien, 315 einen laengeren Volltext, 430 Medien. Ausserdem landen die rohen Antwortbytes gezippt in feed_abrufe, dedupliziert ueber ihren sha256. Der Parser wird sich weiter aendern; nur mit den Rohbytes laesst sich eine Verbesserung rueckwirkend auf die Historie anwenden - dieselbe Ueberlegung wie bei der Trennung von raw_items und codings. Kompression rund 22 %, etwa 20 kB je Abruf. Items ohne Titel werden weiter uebersprungen, jetzt aber gezaehlt. Bestandszeilen bekommen die neuen Felder beim naechsten Lauf nachgetragen, ohne als Revision zu zaehlen; der Inhalts-Hash bleibt dafuer auf Titel/Teaser/URL beschraenkt. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011ZAdZb3EifGb9nvQ5Zp2TE
652 lines
25 KiB
Python
Executable File
652 lines
25 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 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 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)
|
|
dauer = int((time.monotonic() - t0) * 1000)
|
|
senke.notiere(q["id"], slot, dauer, status, gesehen, neu, geaendert,
|
|
uebersprungen, bytes_roh, 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 "
|
|
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())
|