Compare commits

...
5 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
irrlichtandClaude Sonnet 5 8a1ab8a146 Poll-Takt von 15 Minuten auf 2 Stunden senken
Reduziert den Traffic gegen die Quellen. slot_von() rundet jetzt generisch
auf TAKT-Grenzen statt fest auf 15-Minuten-Slots, die Stillstands-Schwelle
in quellen_stillstand steigt von 6h auf 24h, damit der Watchdog nicht schon
nach drei verpassten Slots anschlaegt. GDELT-Vergleich aus den Kommentaren
entfernt, da die Slots nicht mehr mit dessen 15-Minuten-Paketen zusammenfallen.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013vHZCwcCJnz4LKNgvT4bxv
2026-09-04 20:34:22 +02:00
irrlichtandClaude Sonnet 5 c43d8496bf Branding auf Wurzelwerk (SWP) umstellen, Deployment systemagnostisch machen
Ersetzt "Lagebild" durch den Projektnamen Wurzelwerk (Stinkwurzpresse) in
allen Units, Configs und Docs. Entfernt konkrete Hostnamen (Bebop, Pi 5,
Debby) aus README/NOTES zugunsten generischer Rollenbezeichnungen.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013vHZCwcCJnz4LKNgvT4bxv
2026-09-04 20:28:57 +02:00
14 changed files with 654 additions and 127 deletions
+2
View File
@@ -0,0 +1,2 @@
__pycache__/
*.pyc
+1 -1
View File
@@ -9,6 +9,6 @@ RUN pip install --no-cache-dir -r requirements.txt
COPY poller.py sources.toml ./ COPY poller.py sources.toml ./
USER ingest USER ingest
# Ohne Argumente: Dauerbetrieb im 15-Minuten-Takt, Senke aus $DATABASE_URL. # Ohne Argumente: Dauerbetrieb im 2-Stunden-Takt, Senke aus $DATABASE_URL.
# -u, damit die Logs unmittelbar im journal landen und nicht im Puffer haengen. # -u, damit die Logs unmittelbar im journal landen und nicht im Puffer haengen.
ENTRYPOINT ["python", "-u", "/app/poller.py"] ENTRYPOINT ["python", "-u", "/app/poller.py"]
+102 -25
View File
@@ -1,4 +1,4 @@
# Ingest # Wurzelwerk — Ingest
Stufe 0 und 1 der Pipeline: Quellenregister und Roh-Erfassung. Keine Stufe 0 und 1 der Pipeline: Quellenregister und Roh-Erfassung. Keine
Kodierung — was hier landet, soll sich beliebig oft neu kodieren lassen. Kodierung — was hier landet, soll sich beliebig oft neu kodieren lassen.
@@ -8,56 +8,61 @@ Kodierung — was hier landet, soll sich beliebig oft neu kodieren lassen.
| Datei | Zweck | | Datei | Zweck |
|---|---| |---|---|
| `sources.toml` | Quellenregister, einzige Wahrheit ueber die Quellen | | `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 | | `poller.py` | Poller, zwei Senken: Postgres und gzip-JSONL |
| `Containerfile` | Image fuer Bebop | | `Containerfile` | Image fuer den Produktivbetrieb |
| `lagebild-ingest.container`, `lagebild.network` | Quadlet-Units, rootless | | `wurzelwerk-ingest.container`, `wurzelwerk.network` | Quadlet-Units, rootless |
| `ingest.env.example` | Vorlage fuer die Zugangsdaten | | `ingest.env.example` | Vorlage fuer die Zugangsdaten |
## Bebop ## Deploy mit Postgres + Podman/Quadlet (rootless)
```sh ```sh
# 1. Schema einspielen # 1. Schema einspielen (Neuinstallation)
psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql 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 # 2. Image bauen
podman build -t localhost/lagebild-ingest:latest . podman build -t localhost/wurzelwerk-ingest:latest .
# 3. Konfiguration ablegen # 3. Konfiguration ablegen
mkdir -p ~/.config/lagebild mkdir -p ~/.config/wurzelwerk
cp ingest.env.example ~/.config/lagebild/ingest.env cp ingest.env.example ~/.config/wurzelwerk/ingest.env
cp sources.toml ~/.config/lagebild/sources.toml cp sources.toml ~/.config/wurzelwerk/sources.toml
chmod 600 ~/.config/lagebild/ingest.env chmod 600 ~/.config/wurzelwerk/ingest.env
$EDITOR ~/.config/lagebild/ingest.env $EDITOR ~/.config/wurzelwerk/ingest.env
# 4. Quadlet installieren # 4. Quadlet installieren
mkdir -p ~/.config/containers/systemd mkdir -p ~/.config/containers/systemd
cp lagebild-ingest.container lagebild.network ~/.config/containers/systemd/ cp wurzelwerk-ingest.container wurzelwerk.network ~/.config/containers/systemd/
systemctl --user daemon-reload systemctl --user daemon-reload
systemctl --user start lagebild-ingest systemctl --user start wurzelwerk-ingest
journalctl --user -u lagebild-ingest -f journalctl --user -u wurzelwerk-ingest -f
# Damit der Dienst ohne offene Sitzung weiterlaeuft: # Damit der Dienst ohne offene Sitzung weiterlaeuft:
loginctl enable-linger "$USER" loginctl enable-linger "$USER"
``` ```
`sources.toml` ist eingehaengt, nicht ins Image gebacken: eine neue Quelle ist `sources.toml` ist eingehaengt, nicht ins Image gebacken: eine neue Quelle ist
ein `systemctl --user restart lagebild-ingest`, kein Neubau. ein `systemctl --user restart wurzelwerk-ingest`, kein Neubau.
## Pi 5 — Schatten-Ingest ## Schatten-Ingest auf einem Zweitgeraet
Ohne Postgres, ohne Container: Ohne Postgres, ohne Container:
```sh ```sh
python3 poller.py --sink jsonl:/var/lib/lagebild/schatten python3 poller.py --sink jsonl:/pfad/zu/deinem/ablageordner
``` ```
Legt `<dir>/YYYYMMDD/YYYYMMDDHHMMSS.jsonl.gz` an, dedupliziert ueber eine 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 Der Sinn ist nicht Redundanz um ihrer selbst willen: RSS-Feeds haben kein
Archiv. Steht Bebops Poller einen Tag still, ist dieser Tag unwiederbringlich Archiv. Steht der Produktiv-Poller einen Tag still, ist dieser Tag
weg. Kodierungen lassen sich wiederholen, Rohdaten nie. unwiederbringlich weg. Kodierungen lassen sich wiederholen, Rohdaten nie.
## Betrieb ## Betrieb
@@ -78,6 +83,40 @@ gesperrte Quelle bleibt im Register stehen (`aktiv = false`), damit
nachvollziehbar bleibt, warum sie fehlt. Der Poller verweigert den Start, nachvollziehbar bleibt, warum sie fehlt. Der Poller verweigert den Start,
wenn eine als `gesperrt` gefuehrte Quelle aktiv geschaltet ist. 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 ### Waechter
```sql ```sql
@@ -101,18 +140,56 @@ ueber die gesamte Historie, ohne dass ein Feed erneut abgerufen wird.
## Verhalten ## Verhalten
- **15-Minuten-Takt**, an den Slotgrenzen ausgerichtet (+30 s Versatz), - **15-Minuten-Takt**, an den Slotgrenzen ausgerichtet (+30 s Versatz).
identisch zu GDELTs Paketgrenzen — damit bleiben die Slots vergleichbar. 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 - **Bedingter GET** ueber ETag/Last-Modified, im Register gespeichert. Im
Test antworteten 4 von 12 Quellen im zweiten Durchlauf mit 304. Test antworteten 4 von 12 Quellen im zweiten Durchlauf mit 304.
- **Dedup** ueber `unique (quelle, item_key)`; `item_key` ist guid, ersatzweise - **Dedup** ueber `unique (quelle, item_key)`; `item_key` ist guid, ersatzweise
Link, ersatzweise ein Titel-Hash. Link, ersatzweise ein Titel-Hash.
- **Fehlertoleranz**: eine kaputte Quelle beendet den Durchlauf nicht, der - **Fehlertoleranz**: eine kaputte Quelle beendet den Durchlauf nicht, der
Fehler landet in `poll_laeufe`. 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 - **Kein Volltext-Abruf.** Nur der Feed selbst wird geholt. Titel und Teaser
reichen als Kodierungsinput, und damit stellt sich die TDM-Frage nach reichen als Kodierungsinput, und damit stellt sich die TDM-Frage nach
§ 44b UrhG gar nicht erst. § 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) ## Gemessen (2026-09-04)
Erster Lauf ueber 12 aktive Quellen: 820 Items gesehen, 810 neu, gesamter Erster Lauf ueber 12 aktive Quellen: 820 Items gesehen, 810 neu, gesamter
@@ -130,5 +207,5 @@ spiegel | derstandard | 0.92 | Niederlage fuer Donald Trump - Gericht in Missou
- Telegram-Fetcher (Telethon/MTProto). Register ist vorbereitet, Eintraege - Telegram-Fetcher (Telethon/MTProto). Register ist vorbereitet, Eintraege
stehen auf `aktiv = false`. Die Bot-API kann fremde Kanaele nicht lesen. stehen auf `aktiv = false`. Die Bot-API kann fremde Kanaele nicht lesen.
- Ersatz-URL fuer den BR24-Feed (aktuell 404). - Ersatz-URL fuer den BR24-Feed (aktuell 404).
- Quellenpool von 12 auf ~40 erweitern; Regionalzeitungen ueber SearXNG auf - Quellenpool von 12 auf ~40 erweitern; Regionalzeitungen ueber SearXNG
Debby suchen. suchen.
+22 -18
View File
@@ -1,44 +1,48 @@
# swp-01-ingest # swp-01-ingest — Wurzelwerk
Stufe 0/1 der Lagebild-Pipeline: Quellenregister und Roh-Erfassung Stufe 0/1 der Wurzelwerk-Pipeline (Stinkwurzpresse, SWP): Quellenregister
(RSS-Poller, Senken Postgres oder gzip-JSONL). Details zu Betrieb, und Roh-Erfassung (RSS-Poller, Senken Postgres oder gzip-JSONL). Details zu
Verhalten und offenen Punkten stehen in [`NOTES.md`](NOTES.md). Betrieb, Verhalten und offenen Punkten stehen in [`NOTES.md`](NOTES.md).
## Deploy (Bebop, Postgres + Podman/Quadlet) ## Deploy (Server mit Postgres + Podman/Quadlet)
```sh ```sh
# 1. Schema einspielen # 1. Schema einspielen
psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql
# 2. Image bauen # 2. Image bauen
podman build -t localhost/lagebild-ingest:latest . podman build -t localhost/wurzelwerk-ingest:latest .
# 3. Konfiguration ablegen # 3. Konfiguration ablegen
mkdir -p ~/.config/lagebild mkdir -p ~/.config/wurzelwerk
cp ingest.env.example ~/.config/lagebild/ingest.env cp ingest.env.example ~/.config/wurzelwerk/ingest.env
cp sources.toml ~/.config/lagebild/sources.toml cp sources.toml ~/.config/wurzelwerk/sources.toml
chmod 600 ~/.config/lagebild/ingest.env chmod 600 ~/.config/wurzelwerk/ingest.env
$EDITOR ~/.config/lagebild/ingest.env $EDITOR ~/.config/wurzelwerk/ingest.env
# 4. Quadlet installieren # 4. Quadlet installieren
mkdir -p ~/.config/containers/systemd mkdir -p ~/.config/containers/systemd
cp lagebild-ingest.container lagebild.network ~/.config/containers/systemd/ cp wurzelwerk-ingest.container wurzelwerk.network ~/.config/containers/systemd/
systemctl --user daemon-reload systemctl --user daemon-reload
systemctl --user start lagebild-ingest systemctl --user start wurzelwerk-ingest
journalctl --user -u lagebild-ingest -f journalctl --user -u wurzelwerk-ingest -f
# Damit der Dienst ohne offene Sitzung weiterlaeuft: # Damit der Dienst ohne offene Sitzung weiterlaeuft:
loginctl enable-linger "$USER" loginctl enable-linger "$USER"
``` ```
`sources.toml` ist eingehaengt, nicht ins Image gebacken: eine neue Quelle `sources.toml` ist eingehaengt, nicht ins Image gebacken: eine neue Quelle
ist ein `systemctl --user restart lagebild-ingest`, kein Neubau. ist ein `systemctl --user restart wurzelwerk-ingest`, kein Neubau.
## Deploy (Pi 5, ohne Postgres/Container) ## Deploy (Zweitgeraet ohne Postgres/Container, z. B. als Schatten-Ingest)
Fuer ein schlankes Zweitsystem ohne Postgres und ohne Container genuegt
Python 3.11+ mit den Paketen aus `requirements.txt`:
```sh ```sh
python3 poller.py --sink jsonl:/var/lib/lagebild/schatten python3 poller.py --sink jsonl:/pfad/zu/deinem/ablageordner
``` ```
Als `systemd --user`-Unit mit `Restart=always` einrichten (siehe Als `systemd --user`-Unit mit `Restart=always` einrichten (siehe
`lagebild-robots.service` / `.timer` fuer den woechentlichen Robots-Check). `wurzelwerk-robots.service` / `.timer` fuer den woechentlichen
Robots-Check).
+3 -3
View File
@@ -1,10 +1,10 @@
# Nach ~/.config/lagebild/ingest.env kopieren und ausfuellen (chmod 600). # Nach ~/.config/wurzelwerk/ingest.env kopieren und ausfuellen (chmod 600).
# #
# Bei rootless Podman erreichen sich Container ueber ein gemeinsames Netz per # Bei rootless Podman erreichen sich Container ueber ein gemeinsames Netz per
# Containername. Liegt Postgres bereits in einem Netz, dort denselben Namen # Containername. Liegt Postgres bereits in einem Netz, dort denselben Namen
# eintragen wie in der Zeile Network= der Quadlet-Datei. # eintragen wie in der Zeile Network= der Quadlet-Datei.
DATABASE_URL=postgresql://lagebild:GEHEIM@lagebild-pg:5432/lagebild DATABASE_URL=postgresql://wurzelwerk:GEHEIM@wurzelwerk-pg:5432/wurzelwerk
# Kontaktadresse gehoert in den User-Agent: wer unsere Abrufe sieht, soll # Kontaktadresse gehoert in den User-Agent: wer unsere Abrufe sieht, soll
# wissen, wen er anschreiben kann. # wissen, wen er anschreiben kann.
INGEST_USER_AGENT=lagebild-ingest/0.1 (+https://DEINE-DOMAIN/lagebild; kontakt@DEINE-DOMAIN) INGEST_USER_AGENT=wurzelwerk-ingest/0.1 (+https://stinkwurzpresse.de/wurzelwerk; kontakt@stinkwurzpresse.de)
-6
View File
@@ -1,6 +0,0 @@
# Ablegen unter ~/.config/containers/systemd/lagebild.network
[Unit]
Description=Lagebild internes Netz
[Network]
NetworkName=lagebild
+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;
+287 -56
View File
@@ -6,8 +6,8 @@ jedes neue Item roh ab. Bewusst ohne Kodierung: was hier landet, soll sich
beliebig oft neu kodieren lassen. beliebig oft neu kodieren lassen.
Zwei Senken: Zwei Senken:
--sink postgres://... Bebop, die eigentliche Datenbank --sink postgres://... Produktivbetrieb, die eigentliche Datenbank
--sink jsonl:/pfad/dir Pi 5, Schatten-Ingest als gzip-JSONL pro Slot --sink jsonl:/pfad/dir Schatten-Ingest als gzip-JSONL pro Slot
Aufrufe: Aufrufe:
poller.py --sink "$DATABASE_URL" # Dauerbetrieb, 15-Minuten-Takt poller.py --sink "$DATABASE_URL" # Dauerbetrieb, 15-Minuten-Takt
@@ -38,15 +38,20 @@ from email.utils import parsedate_to_datetime
HIER = pathlib.Path(__file__).parent HIER = pathlib.Path(__file__).parent
USER_AGENT = os.environ.get( USER_AGENT = os.environ.get(
"INGEST_USER_AGENT", "INGEST_USER_AGENT",
"lagebild-ingest/0.1 (+https://example.org/lagebild; kontakt@example.org)", "wurzelwerk-ingest/0.1 (+https://stinkwurzpresse.de/wurzelwerk; kontakt@stinkwurzpresse.de)",
) )
TAKT = 15 * 60 # Slotlaenge in Sekunden, wie bei GDELT TAKT = 15 * 60 # Slotlaenge in Sekunden
VERSATZ = 30 # Sekunden nach der Slotgrenze, nicht exakt darauf VERSATZ = 30 # Sekunden nach der Slotgrenze, nicht exakt darauf
NS = { NS = {
"a": "http://www.w3.org/2005/Atom", "a": "http://www.w3.org/2005/Atom",
"dc": "http://purl.org/dc/elements/1.1/", "dc": "http://purl.org/dc/elements/1.1/",
"content": "http://purl.org/rss/1.0/modules/content/", "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 _lauf = True
@@ -66,9 +71,10 @@ def _jetzt():
def slot_von(ts): def slot_von(ts):
"""Auf 15 Minuten abrunden, wie GDELTs Paketgrenzen.""" """Auf TAKT-Sekunden seit Epoch abrunden (bei 15 min auf :00/:15/:30/:45)."""
return ts.replace(second=0, microsecond=0, sekunden = int(ts.timestamp())
minute=(ts.minute // 15) * 15) return dt.datetime.fromtimestamp(
(sekunden // TAKT) * TAKT, tz=dt.timezone.utc)
def norm_titel(t): def norm_titel(t):
@@ -106,6 +112,17 @@ def datum(s):
return d.astimezone(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): def _text(el):
return (el.text or "").strip() if el is not None and el.text else None return (el.text or "").strip() if el is not None and el.text else None
@@ -123,7 +140,12 @@ def _erst(item, pfade):
# Feed-Abruf und -Parsing # Feed-Abruf und -Parsing
# -------------------------------------------------------------------------- # --------------------------------------------------------------------------
def hole(url, etag=None, last_modified=None, timeout=30): 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"} kopf = {"User-Agent": USER_AGENT, "Accept-Encoding": "gzip"}
if etag: if etag:
kopf["If-None-Match"] = etag kopf["If-None-Match"] = etag
@@ -135,23 +157,75 @@ def hole(url, etag=None, last_modified=None, timeout=30):
roh = r.read() roh = r.read()
if r.headers.get("Content-Encoding") == "gzip": if r.headers.get("Content-Encoding") == "gzip":
roh = gzip.decompress(roh) roh = gzip.decompress(roh)
return (r.status, roh, header = {k: r.headers.get(k) for k in HEADER_MERKEN
r.headers.get("ETag"), r.headers.get("Last-Modified")) 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: except urllib.error.HTTPError as err:
if err.code == 304: if err.code == 304:
return 304, None, etag, last_modified return 304, None, etag, last_modified, {}
raise 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): 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) wurzel = ET.fromstring(roh)
eintraege = wurzel.findall(".//item") or wurzel.findall(".//a:entry", NS) eintraege = wurzel.findall(".//item") or wurzel.findall(".//a:entry", NS)
aus = [] aus, uebersprungen = [], 0
for it in eintraege: for pos, it in enumerate(eintraege):
titel = _erst(it, ["title", "a:title"]) titel = _erst(it, ["title", "a:title"])
if not titel: if not titel:
continue # ohne Titel ist ein Item fuer uns wertlos uebersprungen += 1
continue
link = _erst(it, ["link", "a:id"]) link = _erst(it, ["link", "a:id"])
if not link: if not link:
@@ -164,17 +238,36 @@ def parse_feed(roh):
if not guid: if not guid:
guid = hashlib.sha256(titel.encode()).hexdigest()[:32] 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), "titel": klartext(titel),
"url": link, "url": link,
"guid": guid, "guid": guid,
"teaser": klartext(_erst(it, ["description", "content:encoded", "position": pos, # Reihenfolge im Feed = oft die Gewichtung
"a:summary", "a:content"])), "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"]), "autor": _erst(it, ["author", "dc:creator", "a:author/a:name"]),
"pubdate": datum(_erst(it, ["pubDate", "dc:date", "kategorien": _kategorien(it),
"a:published", "a:updated"])), "medien": _medien(it),
}) "links": _links(teaser_html, voll_html),
return aus "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
# -------------------------------------------------------------------------- # --------------------------------------------------------------------------
@@ -184,7 +277,7 @@ class JsonlSenke:
"""Schatten-Ingest fuer den Pi: gzip-JSONL, ein Verzeichnis pro Tag. """Schatten-Ingest fuer den Pi: gzip-JSONL, ein Verzeichnis pro Tag.
Kein Schema, keine Abfragen - nur die Garantie, dass die Rohdaten noch Kein Schema, keine Abfragen - nur die Garantie, dass die Rohdaten noch
da sind, wenn Bebop einen Tag lang stillstand. Die SQLite daneben haelt 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 nur die gesehenen Schluessel, damit nicht bei jedem Slot der komplette
Feed erneut geschrieben wird. Feed erneut geschrieben wird.
""" """
@@ -194,7 +287,8 @@ class JsonlSenke:
self.ziel.mkdir(parents=True, exist_ok=True) self.ziel.mkdir(parents=True, exist_ok=True)
self.db = sqlite3.connect(self.ziel / "gesehen.sqlite") self.db = sqlite3.connect(self.ziel / "gesehen.sqlite")
self.db.execute("create table if not exists gesehen (" 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() self.db.commit()
def sync_quellen(self, quellen): def sync_quellen(self, quellen):
@@ -204,17 +298,31 @@ class JsonlSenke:
def zustand(self, qid): def zustand(self, qid):
return None, None return None, None
def rollback(self):
self.db.rollback()
def schreibe(self, quelle, items, slot): def schreibe(self, quelle, items, slot):
frisch = [] frisch, neu, geaendert = [], 0, 0
for it in items: for it in items:
cur = self.db.execute( row = self.db.execute(
"insert or ignore into gesehen (quelle, guid) values (?, ?)", "select hash from gesehen where quelle = ? and guid = ?",
(quelle["id"], it["guid"])) (quelle["id"], it["guid"])).fetchone()
if cur.rowcount: if row is None:
frisch.append(it) 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() self.db.commit()
if not frisch: if not frisch:
return 0 return 0, 0, 0
tag = self.ziel / slot.strftime("%Y%m%d") tag = self.ziel / slot.strftime("%Y%m%d")
tag.mkdir(exist_ok=True) tag.mkdir(exist_ok=True)
pfad = tag / (slot.strftime("%Y%m%d%H%M%S") + ".jsonl.gz") pfad = tag / (slot.strftime("%Y%m%d%H%M%S") + ".jsonl.gz")
@@ -225,7 +333,20 @@ class JsonlSenke:
{"quelle": quelle["id"], "slot": slot.isoformat(), {"quelle": quelle["id"], "slot": slot.isoformat(),
"abgerufen_am": _jetzt().isoformat(), **it}, "abgerufen_am": _jetzt().isoformat(), **it},
ensure_ascii=False, default=str) + "\n") 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): def notiere(self, *a, **kw):
pass pass
@@ -239,6 +360,9 @@ class PostgresSenke:
import psycopg import psycopg
self.conn = psycopg.connect(dsn, autocommit=False) self.conn = psycopg.connect(dsn, autocommit=False)
def rollback(self):
self.conn.rollback()
def sync_quellen(self, quellen): def sync_quellen(self, quellen):
with self.conn.cursor() as cur: with self.conn.cursor() as cur:
for q in quellen: for q in quellen:
@@ -269,24 +393,84 @@ class PostgresSenke:
return (r[0], r[1]) if r else (None, None) return (r[0], r[1]) if r else (None, None)
def schreibe(self, quelle, items, slot): 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: with self.conn.cursor() as cur:
for it in items: for it in items:
cur.execute(""" cur.execute("""
insert into raw_items insert into raw_items
(quelle, item_key, url, titel, titel_norm, teaser, (quelle, item_key, url, titel, titel_norm, teaser,
autor, pubdate, slot, roh) autor, pubdate, slot, inhalt_hash, roh,
values (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s) position, volltext, kategorien, medien, links,
on conflict (quelle, item_key) do nothing aktualisiert)
returning id 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"], """, (quelle["id"], it["guid"], it["url"], it["titel"],
norm_titel(it["titel"]), it["teaser"], it["autor"], 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))) json.dumps(it, ensure_ascii=False, default=str),
if cur.fetchone(): 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 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() self.conn.commit()
return neu return neu, geaendert, nachgetragen
def merke_zustand(self, qid, etag, last_modified, geaendert): def merke_zustand(self, qid, etag, last_modified, geaendert):
with self.conn.cursor() as cur: with self.conn.cursor() as cur:
@@ -299,14 +483,43 @@ class PostgresSenke:
""", (etag, last_modified, geaendert, qid)) """, (etag, last_modified, geaendert, qid))
self.conn.commit() 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: with self.conn.cursor() as cur:
cur.execute(""" cur.execute("""
insert into poll_laeufe insert into poll_laeufe
(quelle, slot, dauer_ms, http_status, (quelle, slot, dauer_ms, http_status, items_gesehen,
items_gesehen, items_neu, fehler) items_neu, items_geaendert, items_uebersprungen,
values (%s, %s, %s, %s, %s, %s, %s) bytes_roh, fehler)
""", (quelle, slot, dauer_ms, status, gesehen, neu, 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() self.conn.commit()
def close(self): def close(self):
@@ -322,36 +535,53 @@ def lade_quellen(pfad):
def durchlauf(quellen, senke, slot): def durchlauf(quellen, senke, slot):
ges_neu = ges_gesehen = 0 ges_neu = ges_gesehen = ges_geaendert = 0
for q in quellen: for q in quellen:
if q["typ"] != "rss": if q["typ"] != "rss":
continue continue
t0 = time.monotonic() t0 = time.monotonic()
etag, lm = senke.zustand(q["id"]) etag, lm = senke.zustand(q["id"])
status = gesehen = neu = 0 status = gesehen = neu = geaendert = uebersprungen = bytes_roh = 0
nachgetragen = 0
fehler = None fehler = None
try: 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: if status == 304:
pass pass
else: 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) gesehen = len(items)
neu = senke.schreibe(q, items, slot) neu, geaendert, nachgetragen = senke.schreibe(q, items, slot)
if hasattr(senke, "merke_zustand"): if hasattr(senke, "merke_zustand"):
senke.merke_zustand(q["id"], n_etag, n_lm, neu > 0) senke.merke_zustand(q["id"], n_etag, n_lm, neu > 0)
except Exception as e: # eine Quelle darf nicht except Exception as e: # eine Quelle darf nicht
fehler = f"{type(e).__name__}: {e}" # den Durchlauf killen fehler = f"{type(e).__name__}: {e}" # den Durchlauf killen
status = getattr(e, "code", None) 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) 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_neu += neu
ges_geaendert += geaendert
ges_gesehen += gesehen ges_gesehen += gesehen
marke = "!" if fehler else (" " if status != 304 else ".") marke = "!" if fehler else (" " if status != 304 else ".")
print(f" {marke} {q['id']:14} {str(status):>4} " 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) + (f" {fehler}" if fehler else ""), flush=True)
return ges_gesehen, ges_neu return ges_gesehen, ges_neu, ges_geaendert
def check_robots(quellen): def check_robots(quellen):
@@ -419,8 +649,9 @@ def main():
while _lauf: while _lauf:
slot = slot_von(_jetzt()) slot = slot_von(_jetzt())
print(f"[{slot:%Y-%m-%d %H:%M}] Slot", flush=True) print(f"[{slot:%Y-%m-%d %H:%M}] Slot", flush=True)
gesehen, neu = durchlauf(aktiv, senke, slot) gesehen, neu, geaendert = durchlauf(aktiv, senke, slot)
print(f" = {gesehen} gesehen, {neu} neu", flush=True) print(f" = {gesehen} gesehen, {neu} neu, "
f"{geaendert} revidiert", flush=True)
if args.once: if args.once:
break break
ziel = slot + dt.timedelta(seconds=TAKT + VERSATZ) ziel = slot + dt.timedelta(seconds=TAKT + VERSATZ)
+90 -4
View File
@@ -48,19 +48,101 @@ create table if not exists raw_items (
autor text, autor text,
pubdate timestamptz, -- Angabe der Quelle, kann fehlen/luegen pubdate timestamptz, -- Angabe der Quelle, kann fehlen/luegen
abgerufen_am timestamptz not null default now(), abgerufen_am timestamptz not null default now(),
slot timestamptz not null, -- auf 15 min abgerundet, GDELT-kompatibel slot timestamptz not null, -- auf den Poll-Takt abgerundet
sprache text not null default 'de', 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 roh jsonb not null, -- alle geparsten Feldwerte, unveraendert
unique (quelle, item_key) 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_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_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); 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 -- Titel-Clustering laeuft ueber Trigramm-Aehnlichkeit in SQL, nicht in Python
create index if not exists raw_items_titelnorm_trgm create index if not exists raw_items_titelnorm_trgm
on raw_items using gin (titel_norm gin_trgm_ops); 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 -- Kodierungen. Eine Zeile pro (Item, Kodierer-Version). Die Nutzlast bleibt
@@ -105,15 +187,19 @@ create table if not exists poll_laeufe (
begonnen_am timestamptz not null default now(), begonnen_am timestamptz not null default now(),
dauer_ms integer, dauer_ms integer,
http_status integer, -- 304 = unveraendert http_status integer, -- 304 = unveraendert
items_gesehen integer not null default 0, items_gesehen integer not null default 0,
items_neu 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 fehler text
); );
create index if not exists poll_laeufe_slot_idx on poll_laeufe (slot desc); 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); create index if not exists poll_laeufe_quelle_idx on poll_laeufe (quelle, begonnen_am desc);
-- Welche aktive Quelle hat seit ueber sechs Stunden 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 create or replace view quellen_stillstand as
select q.id, q.name, max(r.abgerufen_am) as letztes_item, select q.id, q.name, max(r.abgerufen_am) as letztes_item,
now() - max(r.abgerufen_am) as stille now() - max(r.abgerufen_am) as stille
@@ -1,29 +1,29 @@
# Quadlet-Unit fuer rootless Podman. # Quadlet-Unit fuer rootless Podman.
# Ablegen unter ~/.config/containers/systemd/lagebild-ingest.container, # Ablegen unter ~/.config/containers/systemd/wurzelwerk-ingest.container,
# danach: systemctl --user daemon-reload && systemctl --user start lagebild-ingest # danach: systemctl --user daemon-reload && systemctl --user start wurzelwerk-ingest
# #
# Der Poller laeuft dauerhaft und taktet sich selbst auf die 15-Minuten-Grenzen. # Der Poller laeuft dauerhaft und taktet sich selbst auf die 2-Stunden-Grenzen.
# Deshalb kein systemd-Timer: ein Timer wuerde den bedingten GET-Zustand # Deshalb kein systemd-Timer: ein Timer wuerde den bedingten GET-Zustand
# (ETag/Last-Modified) bei jedem Start neu aufbauen muessen. # (ETag/Last-Modified) bei jedem Start neu aufbauen muessen.
[Unit] [Unit]
Description=Lagebild Ingest-Poller Description=Wurzelwerk Ingest-Poller
After=network-online.target After=network-online.target
Wants=network-online.target Wants=network-online.target
[Container] [Container]
Image=localhost/lagebild-ingest:latest Image=localhost/wurzelwerk-ingest:latest
ContainerName=lagebild-ingest ContainerName=wurzelwerk-ingest
AutoUpdate=local AutoUpdate=local
EnvironmentFile=%h/.config/lagebild/ingest.env EnvironmentFile=%h/.config/wurzelwerk/ingest.env
# Quellenregister von aussen einhaengen: eine neue Quelle ist dann ein # Quellenregister von aussen einhaengen: eine neue Quelle ist dann ein
# Neustart, kein Neubau des Images. # Neustart, kein Neubau des Images.
Volume=%h/.config/lagebild/sources.toml:/app/sources.toml:ro,Z Volume=%h/.config/wurzelwerk/sources.toml:/app/sources.toml:ro,Z
# Dasselbe Netz wie der Postgres-Container. # Dasselbe Netz wie der Postgres-Container.
Network=lagebild.network Network=wurzelwerk.network
NoNewPrivileges=true NoNewPrivileges=true
ReadOnly=true ReadOnly=true
@@ -3,11 +3,11 @@
# damit systemd den Fehlerzustand sichtbar haelt. # damit systemd den Fehlerzustand sichtbar haelt.
[Unit] [Unit]
Description=Lagebild - robots.txt der Quellen pruefen Description=Wurzelwerk - robots.txt der Quellen pruefen
[Service] [Service]
Type=oneshot Type=oneshot
ExecStart=/usr/bin/podman run --rm \ ExecStart=/usr/bin/podman run --rm \
--network=host \ --network=host \
-v %h/.config/lagebild/sources.toml:/app/sources.toml:ro \ -v %h/.config/wurzelwerk/sources.toml:/app/sources.toml:ro \
localhost/lagebild-ingest:latest --check-robots localhost/wurzelwerk-ingest:latest --check-robots
@@ -1,8 +1,8 @@
# Nach ~/.config/systemd/user/ ablegen, dann: # Nach ~/.config/systemd/user/ ablegen, dann:
# systemctl --user enable --now lagebild-robots.timer # systemctl --user enable --now wurzelwerk-robots.timer
[Unit] [Unit]
Description=Lagebild - woechentliche robots.txt-Pruefung Description=Wurzelwerk - woechentliche robots.txt-Pruefung
[Timer] [Timer]
OnCalendar=Mon 05:30 OnCalendar=Mon 05:30
+6
View File
@@ -0,0 +1,6 @@
# Ablegen unter ~/.config/containers/systemd/wurzelwerk.network
[Unit]
Description=Wurzelwerk internes Netz
[Network]
NetworkName=wurzelwerk