Compare commits

..
3 Commits
Author SHA1 Message Date
irrlichtandClaude Opus 5 ab20fa5d15 Fehler einer Quelle nicht mehr die Transaktion mitreissen lassen
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
2026-09-06 12:51:13 +02:00
adminandClaude Opus 5 5a9bbe6e92 Rohe Feed-Antworten ablegen und Parser erweitern
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
2026-09-05 18:21:06 +02:00
adminandClaude Opus 5 4983662514 Poll-Takt auf 15 Minuten und Revisionshistorie fuer Items
Der 2-Stunden-Takt verlor ganze Artikel. Ein RSS-Feed ist ein Fenster,
kein Archiv: ist das Fenster kuerzer als der Poll-Abstand, fallen
Eintraege zwischen zwei Abrufen heraus und sind nicht nachholbar.
Gemessen am 2026-09-04 haelt zeit.de 0,6 h vor, spiegel und welt rund
5 h; von der Zeit sah man etwa jeden dritten Artikel. Der bedingte GET
macht haeufiges Pollen billig, ein unveraenderter Abruf ist ein 304
ohne Body.

Zweitens schrieb der Upsert Aenderungen nicht fort. Redaktionen aendern
Ueberschriften nach der Veroeffentlichung; bei gleicher GUID lief das in
ein `on conflict do nothing` und blieb unsichtbar. raw_items haelt jetzt
den aktuellen Stand und zaehlt Revisionen, raw_item_versionen haelt jede
Fassung einschliesslich der ersten. Erkannt wird das ueber einen
Inhalts-Hash aus Titel, Teaser und URL.

Die Migration zieht Bestandsdaten nach. Die SQL-Hashformel ist
zeichengleich zu poller.inhalt_hash(), geprueft an 814 migrierten Items:
null Scheinrevisionen beim ersten Lauf danach.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_011ZAdZb3EifGb9nvQ5Zp2TE
2026-09-05 18:17:07 +02:00
5 changed files with 581 additions and 62 deletions
+82 -5
View File
@@ -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,13 @@ 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, 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 .
@@ -53,7 +58,7 @@ python3 poller.py --sink jsonl:/pfad/zu/deinem/ablageordner
```
Legt `<dir>/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 `<dir>/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
@@ -78,6 +83,40 @@ 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;
```
### 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
@@ -101,18 +140,56 @@ 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
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 (`<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.
- **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
+80
View File
@@ -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;
+47
View File
@@ -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;
+282 -52
View File
@@ -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,13 +40,18 @@ 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",
"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
@@ -66,7 +71,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 +112,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
@@ -124,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
@@ -136,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):
"""<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."""
"""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:
@@ -165,17 +238,36 @@ def parse_feed(roh):
if not guid:
guid = hashlib.sha256(titel.encode()).hexdigest()[:32]
aus.append({
# 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"])),
})
return aus
"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
# --------------------------------------------------------------------------
@@ -195,7 +287,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):
@@ -205,17 +298,31 @@ class JsonlSenke:
def zustand(self, qid):
return None, None
def rollback(self):
self.db.rollback()
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, 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 +333,20 @@ 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, 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
@@ -240,6 +360,9 @@ class PostgresSenke:
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:
@@ -270,24 +393,84 @@ 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 = 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, 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,
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,
json.dumps(it, ensure_ascii=False, default=str)))
if cur.fetchone():
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
return neu, geaendert, nachgetragen
def merke_zustand(self, qid, etag, last_modified, geaendert):
with self.conn.cursor() as cur:
@@ -300,14 +483,43 @@ class PostgresSenke:
""", (etag, last_modified, geaendert, qid))
self.conn.commit()
def notiere(self, quelle, slot, dauer_ms, status, gesehen, neu, fehler):
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, fehler)
values (%s, %s, %s, %s, %s, %s, %s)
""", (quelle, slot, dauer_ms, status, gesehen, neu, fehler))
(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):
@@ -323,36 +535,53 @@ 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 = 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 = 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
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)
senke.notiere(q["id"], slot, dauer, status, gesehen, neu, fehler)
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 {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
return ges_gesehen, ges_neu, ges_geaendert
def check_robots(quellen):
@@ -420,8 +649,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)
+90 -5
View File
@@ -50,17 +50,99 @@ 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[], -- <category> / dc:subject / Atom-term
medien jsonb, -- enclosure und Media RSS
links text[], -- Verweise aus dem Teaser-HTML
aktualisiert timestamptz, -- Atom <updated>, neben pubdate
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.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';
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;
-- --------------------------------------------------------------------------
-- 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
@@ -105,16 +187,19 @@ 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,
items_uebersprungen integer not null default 0, -- ohne Titel, unbrauchbar
bytes_roh integer,
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 +208,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;