diff --git a/NOTES.md b/NOTES.md index e6e593b..ee09c91 100644 --- a/NOTES.md +++ b/NOTES.md @@ -8,7 +8,8 @@ Kodierung — was hier landet, soll sich beliebig oft neu kodieren lassen. | Datei | Zweck | |---|---| | `sources.toml` | Quellenregister, einzige Wahrheit ueber die Quellen | -| `schema.sql` | Postgres-Schema (`quellen`, `raw_items`, `codings`, `poll_laeufe`) | +| `schema.sql` | Postgres-Schema, vollstaendiger Stand fuer Neuinstallationen | +| `migrations/` | Aenderungen fuer bestehende Installationen, aufsteigend anwenden | | `poller.py` | Poller, zwei Senken: Postgres und gzip-JSONL | | `Containerfile` | Image fuer den Produktivbetrieb | | `wurzelwerk-ingest.container`, `wurzelwerk.network` | Quadlet-Units, rootless | @@ -17,9 +18,12 @@ Kodierung — was hier landet, soll sich beliebig oft neu kodieren lassen. ## Deploy mit Postgres + Podman/Quadlet (rootless) ```sh -# 1. Schema einspielen +# 1. Schema einspielen (Neuinstallation) psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql +# Bestehende Installation stattdessen migrieren: +# psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f migrations/001-revisionen.sql + # 2. Image bauen podman build -t localhost/wurzelwerk-ingest:latest . @@ -78,6 +82,20 @@ 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. +### Ueberschriften-Revisionen + +```sql +-- welche Meldungen wurden nachtraeglich umbenannt +select quelle, fassungen, erste_fassung, letzte_fassung, zuletzt +from titel_revisionen order by zuletzt desc limit 20; + +-- welche Redaktion revidiert am haeufigsten +select quelle, count(*) filter (where revisionen > 0) as revidiert, + count(*) as gesamt, + round(100.0 * count(*) filter (where revisionen > 0) / count(*), 1) as quote +from raw_items group by quelle order by quote desc; +``` + ### Waechter ```sql @@ -101,8 +119,19 @@ ueber die gesamte Historie, ohne dass ein Feed erneut abgerufen wird. ## Verhalten -- **2-Stunden-Takt**, an den Slotgrenzen ausgerichtet (+30 s Versatz), um - unnoetigen Traffic bei den Quellen zu vermeiden. +- **15-Minuten-Takt**, an den Slotgrenzen ausgerichtet (+30 s Versatz). + Der Takt ist nicht frei waehlbar: ein RSS-Feed ist ein Fenster, kein Archiv. + Ist das Fenster kuerzer als der Poll-Abstand, fallen Artikel zwischen zwei + Abrufen heraus und sind unwiederbringlich weg. Gemessen am 2026-09-04: + zeit.de haelt 0,6 h vor, dlf gar nicht messbar, spiegel und welt rund 5 h -- + bei 2 h Takt sah man von der Zeit etwa jeden dritten Artikel. Der bedingte + GET macht haeufiges Pollen billig: ein unveraenderter Abruf ist ein 304 ohne + Body. +- **Revisionen** werden mitgeschrieben. Redaktionen aendern Ueberschriften nach + der Veroeffentlichung; frueher lief das in ein `on conflict do nothing` und + blieb unsichtbar. `raw_items` haelt den aktuellen Stand, `raw_item_versionen` + jede Fassung. Der Datensatz existiert sonst nirgends -- GDELT sieht jede URL + genau einmal. - **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 diff --git a/migrations/001-revisionen.sql b/migrations/001-revisionen.sql new file mode 100644 index 0000000..e279b49 --- /dev/null +++ b/migrations/001-revisionen.sql @@ -0,0 +1,80 @@ +-- Migration 001: Revisionshistorie fuer Items +-- +-- Fuer bestehende Installationen. Bei einer Neuinstallation deckt schema.sql +-- das bereits ab; diese Datei ist dann ein No-op. +-- +-- psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f migrations/001-revisionen.sql + +begin; + +alter table raw_items + add column if not exists inhalt_hash text, + add column if not exists revisionen integer not null default 0, + add column if not exists zuletzt_geaendert timestamptz; + +-- Bestandsdaten nachziehen. Die Formel muss zeichengleich zu +-- poller.inhalt_hash() sein, sonst gilt beim naechsten Lauf jedes Altitem +-- als geaendert. Trennzeichen ist U+001F (unit separator). +update raw_items + set inhalt_hash = md5(coalesce(titel, '') || E'\x1f' || + coalesce(teaser, '') || E'\x1f' || + coalesce(url, '')) + where inhalt_hash is null; + +alter table raw_items alter column inhalt_hash set not null; + +create index if not exists raw_items_revidiert_idx + on raw_items (zuletzt_geaendert desc) where revisionen > 0; + +create table if not exists raw_item_versionen ( + id bigserial primary key, + raw_item_id bigint not null references raw_items(id) on delete cascade, + gesehen_am timestamptz not null default now(), + slot timestamptz not null, + titel text not null, + teaser text, + url text, + inhalt_hash text not null, + + unique (raw_item_id, inhalt_hash) +); + +create index if not exists raw_item_versionen_item_idx + on raw_item_versionen (raw_item_id, gesehen_am); + +-- Erstfassung der Bestandsitems eintragen, damit die Historie vollstaendig ist. +-- gesehen_am auf abgerufen_am setzen, nicht auf now(): die Fassung ist alt. +insert into raw_item_versionen (raw_item_id, gesehen_am, slot, titel, teaser, + url, inhalt_hash) +select id, abgerufen_am, slot, titel, teaser, url, inhalt_hash +from raw_items +on conflict (raw_item_id, inhalt_hash) do nothing; + +alter table poll_laeufe + add column if not exists items_geaendert integer not null default 0; + +drop view if exists titel_revisionen; +create view titel_revisionen as +select r.id, r.quelle, r.url, r.revisionen, + count(*) as fassungen, + (array_agg(v.titel order by v.gesehen_am))[1] as erste_fassung, + (array_agg(v.titel order by v.gesehen_am desc))[1] as letzte_fassung, + min(v.gesehen_am) as erstmals, + max(v.gesehen_am) as zuletzt +from raw_items r +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; + +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; + +commit; diff --git a/poller.py b/poller.py index f23bc95..1cb35df 100755 --- a/poller.py +++ b/poller.py @@ -1,7 +1,7 @@ #!/usr/bin/env python3 """Ingest-Poller fuer den deutschen Ereignisdatensatz. -Liest die Quellen aus sources.toml, holt alle 2 Stunden die Feeds und legt +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. @@ -10,7 +10,7 @@ Zwei Senken: --sink jsonl:/pfad/dir Schatten-Ingest als gzip-JSONL pro Slot Aufrufe: - poller.py --sink "$DATABASE_URL" # Dauerbetrieb, 2-Stunden-Takt + poller.py --sink "$DATABASE_URL" # Dauerbetrieb, 15-Minuten-Takt poller.py --sink "$DATABASE_URL" --once # ein Durchlauf poller.py --check-robots # Lizenzstatus nachpruefen """ @@ -40,7 +40,7 @@ USER_AGENT = os.environ.get( "INGEST_USER_AGENT", "wurzelwerk-ingest/0.1 (+https://stinkwurzpresse.de/wurzelwerk; kontakt@stinkwurzpresse.de)", ) -TAKT = 2 * 60 * 60 # Slotlaenge in Sekunden +TAKT = 15 * 60 # Slotlaenge in Sekunden VERSATZ = 30 # Sekunden nach der Slotgrenze, nicht exakt darauf NS = { "a": "http://www.w3.org/2005/Atom", @@ -66,7 +66,7 @@ def _jetzt(): def slot_von(ts): - """Auf TAKT-Sekunden seit Epoch abrunden (bei 2 h also volle gerade Stunden).""" + """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) @@ -107,6 +107,17 @@ def datum(s): 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 @@ -165,7 +176,7 @@ def parse_feed(roh): if not guid: guid = hashlib.sha256(titel.encode()).hexdigest()[:32] - aus.append({ + eintrag = { "titel": klartext(titel), "url": link, "guid": guid, @@ -174,7 +185,10 @@ def parse_feed(roh): "autor": _erst(it, ["author", "dc:creator", "a:author/a:name"]), "pubdate": datum(_erst(it, ["pubDate", "dc:date", "a:published", "a:updated"])), - }) + } + eintrag["inhalt_hash"] = inhalt_hash( + eintrag["titel"], eintrag["teaser"], eintrag["url"]) + aus.append(eintrag) return aus @@ -195,7 +209,8 @@ class JsonlSenke: 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))") + "quelle text, guid text, hash text, " + "primary key (quelle, guid))") self.db.commit() def sync_quellen(self, quellen): @@ -206,16 +221,27 @@ class JsonlSenke: return None, None def schreibe(self, quelle, items, slot): - frisch = [] + frisch, neu, geaendert = [], 0, 0 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) + 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 + return 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") @@ -226,7 +252,7 @@ class JsonlSenke: {"quelle": quelle["id"], "slot": slot.isoformat(), "abgerufen_am": _jetzt().isoformat(), **it}, ensure_ascii=False, default=str) + "\n") - return len(frisch) + return neu, geaendert def notiere(self, *a, **kw): pass @@ -270,24 +296,57 @@ class PostgresSenke: return (r[0], r[1]) if r else (None, None) def schreibe(self, quelle, items, slot): - neu = 0 + """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 = 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 + autor, pubdate, slot, inhalt_hash, roh) + values (%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, + inhalt_hash = excluded.inhalt_hash, + roh = excluded.roh, + revisionen = raw_items.revisionen + 1, + zuletzt_geaendert = now() + where raw_items.inhalt_hash is distinct from excluded.inhalt_hash + 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["pubdate"], slot, it["inhalt_hash"], json.dumps(it, ensure_ascii=False, default=str))) - if cur.fetchone(): + zeile = cur.fetchone() + if not zeile: + continue # unveraendert + rid, ist_neu = zeile + if ist_neu: neu += 1 + else: + geaendert += 1 + 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 + """, (rid, slot, it["titel"], it["teaser"], it["url"], + it["inhalt_hash"])) self.conn.commit() - return neu + return neu, geaendert def merke_zustand(self, qid, etag, last_modified, geaendert): with self.conn.cursor() as cur: @@ -300,14 +359,16 @@ class PostgresSenke: """, (etag, last_modified, geaendert, qid)) self.conn.commit() - def notiere(self, quelle, slot, dauer_ms, status, gesehen, neu, fehler): + def notiere(self, quelle, slot, dauer_ms, status, gesehen, neu, + geaendert, 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)) + items_gesehen, items_neu, items_geaendert, fehler) + values (%s, %s, %s, %s, %s, %s, %s, %s) + """, (quelle, slot, dauer_ms, status, gesehen, neu, geaendert, + fehler)) self.conn.commit() def close(self): @@ -323,13 +384,13 @@ def lade_quellen(pfad): def durchlauf(quellen, senke, slot): - ges_neu = ges_gesehen = 0 + 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 = 0 + status = gesehen = neu = geaendert = 0 fehler = None try: status, roh, n_etag, n_lm = hole(q["url"], etag, lm) @@ -338,21 +399,23 @@ def durchlauf(quellen, senke, slot): else: items = parse_feed(roh) gesehen = len(items) - neu = senke.schreibe(q, items, slot) + neu, geaendert = 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) + senke.notiere(q["id"], slot, dauer, status, gesehen, neu, geaendert, + 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 {dauer:5} ms" + f"{gesehen:4} gesehen {neu:4} neu {geaendert:3} rev {dauer:5} ms" + (f" {fehler}" if fehler else ""), flush=True) - return ges_gesehen, ges_neu + return ges_gesehen, ges_neu, ges_geaendert def check_robots(quellen): @@ -420,8 +483,9 @@ def main(): 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) + 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) diff --git a/schema.sql b/schema.sql index a0382dc..6f4ea66 100644 --- a/schema.sql +++ b/schema.sql @@ -50,17 +50,63 @@ 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', + inhalt_hash text not null, -- md5 ueber titel|teaser|url + revisionen integer not null default 0, + zuletzt_geaendert timestamptz, roh jsonb not null, -- alle geparsten Feldwerte, unveraendert unique (quelle, item_key) ); +comment on column raw_items.revisionen is + 'Wie oft die Quelle das Item nach der Erstveroeffentlichung geaendert hat'; + 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); +create index if not exists raw_items_revidiert_idx + on raw_items (zuletzt_geaendert desc) where revisionen > 0; + +-- -------------------------------------------------------------------------- +-- Fassungen. Redaktionen aendern Ueberschriften nach der Veroeffentlichung; +-- frueher lief das in ein `on conflict do nothing` und blieb unsichtbar. +-- raw_items haelt den aktuellen Stand, diese Tabelle jede Fassung - auch die +-- erste, damit die Historie vollstaendig ist. +-- +-- Der Datensatz existiert sonst nirgends: GDELT sieht jede URL genau einmal. +-- -------------------------------------------------------------------------- +create table if not exists raw_item_versionen ( + id bigserial primary key, + raw_item_id bigint not null references raw_items(id) on delete cascade, + gesehen_am timestamptz not null default now(), + slot timestamptz not null, + titel text not null, + teaser text, + url text, + inhalt_hash text not null, + + unique (raw_item_id, inhalt_hash) +); + +create index if not exists raw_item_versionen_item_idx + on raw_item_versionen (raw_item_id, gesehen_am); + +-- Items, deren Ueberschrift sich geaendert hat, mit erster und letzter Fassung +drop view if exists titel_revisionen; +create view titel_revisionen as +select r.id, r.quelle, r.url, r.revisionen, + count(*) as fassungen, + (array_agg(v.titel order by v.gesehen_am))[1] as erste_fassung, + (array_agg(v.titel order by v.gesehen_am desc))[1] as letzte_fassung, + min(v.gesehen_am) as erstmals, + max(v.gesehen_am) as zuletzt +from raw_items r +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; -- -------------------------------------------------------------------------- -- Kodierungen. Eine Zeile pro (Item, Kodierer-Version). Die Nutzlast bleibt @@ -105,16 +151,17 @@ create table if not exists poll_laeufe ( 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, + items_gesehen integer not null default 0, + items_neu integer not null default 0, + items_geaendert 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 24 Stunden (12 verpasste Slots bei --- 2-Stunden-Takt) nichts Neues geliefert? +-- Welche aktive Quelle hat seit ueber sechs Stunden (24 verpasste Slots bei +-- 15-Minuten-Takt) 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 @@ -123,5 +170,5 @@ 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 '24 hours' + or now() - max(r.abgerufen_am) > interval '6 hours' order by stille desc nulls first;