diff --git a/NOTES.md b/NOTES.md index ee09c91..1c9003b 100644 --- a/NOTES.md +++ b/NOTES.md @@ -21,8 +21,9 @@ Kodierung — was hier landet, soll sich beliebig oft neu kodieren lassen. # 1. Schema einspielen (Neuinstallation) psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql -# Bestehende Installation stattdessen migrieren: +# Bestehende Installation stattdessen migrieren, aufsteigend: # psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f migrations/001-revisionen.sql +# psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f migrations/002-rohablage-und-felder.sql # 2. Image bauen podman build -t localhost/wurzelwerk-ingest:latest . @@ -57,7 +58,7 @@ python3 poller.py --sink jsonl:/pfad/zu/deinem/ablageordner ``` Legt `/YYYYMMDD/YYYYMMDDHHMMSS.jsonl.gz` an, dedupliziert ueber eine -SQLite daneben. Als `systemd --user`-Unit mit `Restart=always` einrichten. +SQLite daneben, und die rohen Feeds unter `/YYYYMMDD/feeds/`. Als `systemd --user`-Unit mit `Restart=always` einrichten. Der Sinn ist nicht Redundanz um ihrer selbst willen: RSS-Feeds haben kein Archiv. Steht der Produktiv-Poller einen Tag still, ist dieser Tag @@ -96,6 +97,26 @@ select quelle, count(*) filter (where revisionen > 0) as revidiert, from raw_items group by quelle order by quote desc; ``` +### Was die Verlage selbst labeln + +```sql +-- Themenverteilung je Quelle, ohne einen Kodierer bemueht zu haben +select quelle, k as kategorie, count(*) +from raw_items, unnest(kategorien) k +group by quelle, k order by count(*) desc limit 20; +``` + +### Speicherbedarf im Blick behalten + +```sql +select pg_size_pretty(pg_total_relation_size('feed_abrufe')) as rohablage, + pg_size_pretty(pg_total_relation_size('raw_items')) as items; + +-- Rohablage aelter als ein Jahr wegwerfen, falls es eng wird. raw_items +-- bleibt davon unberuehrt; nur die Moeglichkeit, neu zu parsen, entfaellt. +-- delete from feed_abrufe where abgerufen_am < now() - interval '1 year'; +``` + ### Waechter ```sql @@ -138,10 +159,37 @@ ueber die gesamte Historie, ohne dass ein Feed erneut abgerufen wird. Link, ersatzweise ein Titel-Hash. - **Fehlertoleranz**: eine kaputte Quelle beendet den Durchlauf nicht, der Fehler landet in `poll_laeufe`. +- **Rohe Feed-Antworten** landen gezippt in `feed_abrufe`, dedupliziert ueber + den sha256 der Rohbytes. Der Parser wird sich aendern; nur mit den Rohbytes + laesst sich eine spaetere Verbesserung rueckwirkend auf die Historie + anwenden. Kompression rund 22 %, etwa 20 kB je abgelegtem Abruf. +- **Der Parser nimmt mit, was im XML steht**: Kategorien (``, + `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. +- **Items ohne Titel** werden weiterhin uebersprungen, aber gezaehlt + (`poll_laeufe.items_uebersprungen`) statt stillschweigend verworfen. - **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. +## Migration bestehender Installationen + +`001` und `002` sind idempotent und koennen im laufenden Betrieb eingespielt +werden; danach den Poller neu starten. + +Nach `002` sind die neuen Spalten auf Bestandszeilen zunaechst leer. Der +naechste Lauf traegt sie fuer alle Items nach, die die Quelle noch ausliefert +(in der Ausgabe als `nachgetr`), ohne sie als Revision zu zaehlen. Items, die +inzwischen aus dem Feed gefallen sind, behalten leere Zusatzfelder - die Daten +sind nicht mehr abrufbar. + +Der Inhalts-Hash bleibt bewusst auf Titel, Teaser und URL beschraenkt. Naehme +er die neuen Felder auf, gaelte beim ersten Lauf nach dem Update jedes +Bestandsitem als revidiert. + ## Gemessen (2026-09-04) Erster Lauf ueber 12 aktive Quellen: 820 Items gesehen, 810 neu, gesamter diff --git a/migrations/002-rohablage-und-felder.sql b/migrations/002-rohablage-und-felder.sql new file mode 100644 index 0000000..ec2ebeb --- /dev/null +++ b/migrations/002-rohablage-und-felder.sql @@ -0,0 +1,47 @@ +-- Migration 002: Rohe Feed-Antworten ablegen, Parser-Felder ergaenzen +-- +-- psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f migrations/002-rohablage-und-felder.sql +-- +-- Setzt 001 voraus. Bestandszeilen bekommen die neuen Spalten leer; sie +-- fuellen sich, sobald die Quelle das Item das naechste Mal ausliefert. +-- Der Inhalts-Hash bleibt bewusst auf Titel/Teaser/URL beschraenkt, sonst +-- gaelte beim ersten Lauf nach dem Update jedes Bestandsitem als revidiert. + +begin; + +alter table raw_items + add column if not exists position integer, + add column if not exists volltext text, + add column if not exists kategorien text[], + add column if not exists medien jsonb, + add column if not exists links text[], + add column if not exists aktualisiert timestamptz; + +comment on column raw_items.kategorien is + 'Vom Verlag selbst vergebene Themenlabel - Vorfilter und Prompt-Kontext'; +comment on column raw_items.position is + 'Rang im Feed beim ersten Sehen. Spaeter nicht fortgeschrieben'; + +create table if not exists feed_abrufe ( + id bigserial primary key, + quelle text not null references quellen(id), + slot timestamptz not null, + abgerufen_am timestamptz not null default now(), + http_status integer, + header jsonb, + koerper_hash text not null, + koerper bytea not null, + bytes_roh integer, + bytes_gz integer, + + unique (quelle, koerper_hash) +); + +create index if not exists feed_abrufe_quelle_idx + on feed_abrufe (quelle, abgerufen_am desc); + +alter table poll_laeufe + add column if not exists items_uebersprungen integer not null default 0, + add column if not exists bytes_roh integer; + +commit; diff --git a/poller.py b/poller.py index 1cb35df..952739f 100755 --- a/poller.py +++ b/poller.py @@ -46,7 +46,12 @@ 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 @@ -135,7 +140,12 @@ def _erst(item, pfade): # Feed-Abruf und -Parsing # -------------------------------------------------------------------------- def hole(url, etag=None, last_modified=None, timeout=30): - """Bedingter GET. Gibt (status, bytes|None, etag, last_modified) zurueck.""" + """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 @@ -147,23 +157,75 @@ def hole(url, etag=None, last_modified=None, timeout=30): 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")) + 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 + return 304, None, etag, last_modified, {} raise +_HREF_RE = re.compile(r'href=["\']([^"\']+)["\']', re.I) + + +def _kategorien(it): + """, 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.""" + """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 = [] - for it in eintraege: + aus, uebersprungen = [], 0 + for pos, it in enumerate(eintraege): titel = _erst(it, ["title", "a:title"]) if not titel: - continue # ohne Titel ist ein Item fuer uns wertlos + uebersprungen += 1 + continue link = _erst(it, ["link", "a:id"]) if not link: @@ -176,20 +238,36 @@ def parse_feed(roh): 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, - "teaser": klartext(_erst(it, ["description", "content:encoded", - "a:summary", "a:content"])), + "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"]), - "pubdate": datum(_erst(it, ["pubDate", "dc:date", - "a:published", "a:updated"])), + "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 + return aus, uebersprungen # -------------------------------------------------------------------------- @@ -241,7 +319,7 @@ class JsonlSenke: frisch.append(it) self.db.commit() if not frisch: - return 0, 0 + 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") @@ -252,7 +330,20 @@ class JsonlSenke: {"quelle": quelle["id"], "slot": slot.isoformat(), "abgerufen_am": _jetzt().isoformat(), **it}, ensure_ascii=False, default=str) + "\n") - return neu, geaendert + 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 @@ -305,14 +396,17 @@ class PostgresSenke: ist. Bleibt der Hash gleich, liefert das WHERE keine Zeile zurueck und es passiert nichts. """ - neu = geaendert = 0 + 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) - values (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) + 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, @@ -320,33 +414,57 @@ class PostgresSenke: 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, - revisionen = raw_items.revisionen + 1, - zuletzt_geaendert = now() + 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))) + 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 + continue # unveraendert, nichts zu tun rid, ist_neu = zeile - if ist_neu: - neu += 1 - else: - geaendert += 1 + + # 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 + return neu, geaendert, nachgetragen def merke_zustand(self, qid, etag, last_modified, geaendert): with self.conn.cursor() as cur: @@ -359,16 +477,43 @@ class PostgresSenke: """, (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, fehler): + 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, fehler) - values (%s, %s, %s, %s, %s, %s, %s, %s) + (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, - fehler)) + uebersprungen, bytes_roh, fehler)) self.conn.commit() def close(self): @@ -390,16 +535,19 @@ def durchlauf(quellen, senke, slot): continue t0 = time.monotonic() etag, lm = senke.zustand(q["id"]) - status = gesehen = neu = geaendert = 0 + status = gesehen = neu = geaendert = uebersprungen = bytes_roh = 0 + nachgetragen = 0 fehler = None try: - status, roh, n_etag, n_lm = hole(q["url"], etag, lm) + status, roh, n_etag, n_lm, header = hole(q["url"], etag, lm) if status == 304: pass else: - items = parse_feed(roh) + bytes_roh = len(roh) + senke.speichere_feed(q, slot, status, header, roh) + items, uebersprungen = parse_feed(roh) gesehen = len(items) - neu, geaendert = senke.schreibe(q, items, slot) + 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 @@ -407,13 +555,16 @@ def durchlauf(quellen, senke, slot): status = getattr(e, "code", None) dauer = int((time.monotonic() - t0) * 1000) senke.notiere(q["id"], slot, dauer, status, gesehen, neu, geaendert, - fehler) + 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 {dauer:5} ms" + 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 diff --git a/schema.sql b/schema.sql index 6f4ea66..07a20b7 100644 --- a/schema.sql +++ b/schema.sql @@ -50,6 +50,12 @@ create table if not exists raw_items ( abgerufen_am timestamptz not null default now(), slot timestamptz not null, -- auf den Poll-Takt abgerundet sprache text not null default 'de', + position integer, -- Rang beim ersten Sehen + volltext text, -- content:encoded, falls zusaetzlich da + kategorien text[], -- / dc:subject / Atom-term + medien jsonb, -- enclosure und Media RSS + links text[], -- Verweise aus dem Teaser-HTML + aktualisiert timestamptz, -- Atom , neben pubdate inhalt_hash text not null, -- md5 ueber titel|teaser|url revisionen integer not null default 0, zuletzt_geaendert timestamptz, @@ -58,6 +64,10 @@ create table if not exists raw_items ( unique (quelle, item_key) ); +comment on column raw_items.kategorien is + 'Vom Verlag selbst vergebene Themenlabel - Vorfilter und Prompt-Kontext'; +comment on column raw_items.position is + 'Rang im Feed beim ersten Sehen. Spaeter nicht fortgeschrieben'; comment on column raw_items.revisionen is 'Wie oft die Quelle das Item nach der Erstveroeffentlichung geaendert hat'; @@ -108,6 +118,32 @@ join raw_item_versionen v on v.raw_item_id = r.id where r.revisionen > 0 group by r.id, r.quelle, r.url, r.revisionen; +-- -------------------------------------------------------------------------- +-- Rohe Feed-Antworten. Der Parser wird sich aendern; nur mit den Rohbytes +-- laesst sich eine spaetere Verbesserung rueckwirkend auf die Historie +-- anwenden. Dieselbe Ueberlegung wie bei der Trennung raw_items / codings. +-- +-- Dedupliziert ueber (quelle, koerper_hash): ein unveraenderter Feed wird +-- nicht zweimal abgelegt, und ein 304 hat ohnehin keinen Koerper. +-- -------------------------------------------------------------------------- +create table if not exists feed_abrufe ( + id bigserial primary key, + quelle text not null references quellen(id), + slot timestamptz not null, + abgerufen_am timestamptz not null default now(), + http_status integer, + header jsonb, -- Date, Age, Cache-Control, ... + koerper_hash text not null, -- sha256 ueber die Rohbytes + koerper bytea not null, -- gzip + bytes_roh integer, + bytes_gz integer, + + unique (quelle, koerper_hash) +); + +create index if not exists feed_abrufe_quelle_idx + on feed_abrufe (quelle, abgerufen_am desc); + -- -------------------------------------------------------------------------- -- Kodierungen. Eine Zeile pro (Item, Kodierer-Version). Die Nutzlast bleibt -- absichtlich jsonb, solange nicht gemessen ist, bis zu welcher CAMEO-Ebene @@ -154,6 +190,8 @@ create table if not exists poll_laeufe ( items_gesehen integer not null default 0, items_neu integer not null default 0, items_geaendert integer not null default 0, + items_uebersprungen integer not null default 0, -- ohne Titel, unbrauchbar + bytes_roh integer, fehler text );