Plugin neuheiten: Web-Quellen-Crawler (4 Adapter, Rate-Limit, Dedup-Gate, Sync-UI)
This commit is contained in:
@@ -24,6 +24,7 @@ Keine Fachlogik im Kern — alles hier im Plugin gemäß Plugin-Vertrag.
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import threading
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from urllib.parse import quote
|
||||
@@ -37,7 +38,13 @@ from redaktionskern.auth.models import Role, User
|
||||
from redaktionskern.contracts import BasePlugin, Migration, NavEntry, PluginContext
|
||||
|
||||
from .bgg import BggClient, BggFehler
|
||||
from .models import Neuheit
|
||||
from .models import Neuheit, QuellenStatus
|
||||
from .quellen import ADAPTER_KLASSEN
|
||||
from .quellen_sync import (
|
||||
QuellenLaufZeile,
|
||||
WebQuellenSyncService,
|
||||
quelle_aktiv,
|
||||
)
|
||||
from .sync import SyncErgebnis, SyncService
|
||||
|
||||
_logger = logging.getLogger("plugins.neuheiten")
|
||||
@@ -77,6 +84,7 @@ class NeuheitenPlugin(BasePlugin):
|
||||
super().__init__()
|
||||
self._scheduler = None # BackgroundScheduler, falls aktiviert
|
||||
self._sync_sperre = threading.Lock()
|
||||
self._quellen_sperre = threading.Lock()
|
||||
self._suchbegriffe: list[str] = []
|
||||
self._max_treffer_pro_suche = 25
|
||||
|
||||
@@ -117,6 +125,7 @@ class NeuheitenPlugin(BasePlugin):
|
||||
)
|
||||
]
|
||||
|
||||
quellen_zeilen = self._quellen_uebersicht(db)
|
||||
kontext = {
|
||||
"user": user,
|
||||
"titel": self.title,
|
||||
@@ -135,6 +144,11 @@ class NeuheitenPlugin(BasePlugin):
|
||||
"suchbegriffe": ", ".join(self._suchbegriffe) or "—",
|
||||
"sync_aktiv": self._scheduler is not None and self._scheduler.running,
|
||||
"anzahl": len(eintraege),
|
||||
"quellen_zeilen": quellen_zeilen,
|
||||
"quellen_sync_aktiv": (
|
||||
os.environ.get("SPIELE_NEUHEITEN_QUELLEN_SYNC_AKTIV", "1").strip()
|
||||
== "1"
|
||||
),
|
||||
}
|
||||
if request.headers.get("HX-Request") == "true":
|
||||
return self.context.templates.TemplateResponse(
|
||||
@@ -167,13 +181,85 @@ class NeuheitenPlugin(BasePlugin):
|
||||
ziel = f"/neuheiten?meldung={quote(f'Sync abgeschlossen: {ergebnis.als_text()}')}"
|
||||
return RedirectResponse(ziel, status_code=303)
|
||||
|
||||
@self.router.post("/neuheiten/quellen-sync")
|
||||
def quellen_jetzt_synchronisieren(
|
||||
user: User = Depends(require_roles(Role.ADMIN.value, Role.REDAKTEUR.value)),
|
||||
):
|
||||
zeilen = self._quellen_sync_ausfuehren()
|
||||
text = " · ".join(z.als_text() for z in zeilen)
|
||||
ziel = f"/neuheiten?meldung={quote(f'Web-Quellen-Sync abgeschlossen: {text}')}"
|
||||
return RedirectResponse(ziel[:2000], status_code=303)
|
||||
|
||||
def _quellen_uebersicht(self, db) -> list[dict]:
|
||||
"""Zeilen der Sync-Übersicht: Adapter + letzter Lauf pro Quelle."""
|
||||
status_nach_quelle = {
|
||||
zeile.quelle: zeile
|
||||
for zeile in db.scalars(select(QuellenStatus)).all()
|
||||
}
|
||||
ansichten = []
|
||||
for klasse in ADAPTER_KLASSEN:
|
||||
status = status_nach_quelle.get(klasse.name)
|
||||
ansichten.append(
|
||||
{
|
||||
"name": klasse.name,
|
||||
"anzeigename": klasse.anzeigename,
|
||||
"url": klasse.start_urls[0] if klasse.start_urls else "",
|
||||
"aktiv": quelle_aktiv(klasse.name),
|
||||
"letzte_laufzeit": status.letzte_laufzeit if status else None,
|
||||
"dauer_sekunden": status.dauer_sekunden if status else None,
|
||||
"neu": status.anzahl_neu if status else 0,
|
||||
"aktualisiert": status.anzahl_aktualisiert if status else 0,
|
||||
"gefiltert": status.anzahl_gefiltert if status else 0,
|
||||
"fehleranzahl": status.anzahl_fehler if status else 0,
|
||||
"fehlermeldung": status.fehlermeldung if status else None,
|
||||
"unveraendert": status.unveraendert if status else False,
|
||||
}
|
||||
)
|
||||
return ansichten
|
||||
|
||||
# ---------- Plugin-Vertrag ----------
|
||||
|
||||
def migrations(self) -> list[Migration]:
|
||||
def neuheiten_tabelle(conn) -> None:
|
||||
Neuheit.__table__.create(conn, checkfirst=True)
|
||||
|
||||
return [Migration(version="0001_neuheiten_tabelle", up=neuheiten_tabelle)]
|
||||
def bgg_id_nullable(conn) -> None:
|
||||
"""Migration 0002: Web-Quellen-Einträge ohne BGG-ID erlauben."""
|
||||
if conn.dialect.name == "sqlite":
|
||||
info = conn.exec_driver_sql("PRAGMA table_info(neuheiten)").fetchall()
|
||||
bgg_spalte = next((z for z in info if z[1] == "bgg_id"), None)
|
||||
if bgg_spalte is None or not bgg_spalte[3]:
|
||||
# Spalte fehlt (frisch angelegt) oder ist bereits nullable.
|
||||
return
|
||||
# SQLite kennt kein ALTER COLUMN: Tabelle nach aktuellem Modell
|
||||
# neu anlegen, Daten übernehmen, alte Tabelle verwerfen.
|
||||
conn.exec_driver_sql("ALTER TABLE neuheiten RENAME TO neuheiten_alt_0002")
|
||||
for ix in Neuheit.__table__.indexes:
|
||||
conn.exec_driver_sql(f'DROP INDEX IF EXISTS "{ix.name}"')
|
||||
Neuheit.__table__.create(conn, checkfirst=True)
|
||||
spalten = ", ".join(f'"{c.name}"' for c in Neuheit.__table__.columns)
|
||||
conn.exec_driver_sql(
|
||||
f"INSERT INTO neuheiten ({spalten}) SELECT {spalten} FROM neuheiten_alt_0002"
|
||||
)
|
||||
conn.exec_driver_sql("DROP TABLE neuheiten_alt_0002")
|
||||
else:
|
||||
ist_nullable = conn.exec_driver_sql(
|
||||
"SELECT is_nullable FROM information_schema.columns "
|
||||
"WHERE table_name = 'neuheiten' AND column_name = 'bgg_id'"
|
||||
).scalar()
|
||||
if ist_nullable == "NO":
|
||||
conn.exec_driver_sql(
|
||||
"ALTER TABLE neuheiten ALTER COLUMN bgg_id DROP NOT NULL"
|
||||
)
|
||||
|
||||
def quellen_status_tabelle(conn) -> None:
|
||||
QuellenStatus.__table__.create(conn, checkfirst=True)
|
||||
|
||||
return [
|
||||
Migration(version="0001_neuheiten_tabelle", up=neuheiten_tabelle),
|
||||
Migration(version="0002_bgg_id_nullable", up=bgg_id_nullable),
|
||||
Migration(version="0003_quellen_status", up=quellen_status_tabelle),
|
||||
]
|
||||
|
||||
def navigation(self) -> list[NavEntry]:
|
||||
return [NavEntry(label=self.title, url="/neuheiten")]
|
||||
@@ -192,6 +278,10 @@ class NeuheitenPlugin(BasePlugin):
|
||||
|
||||
if os.environ.get("SPIELE_BGG_SYNC_AKTIV", "1").strip() == "1":
|
||||
self._scheduler_starten(os.environ.get("SPIELE_BGG_SYNC_INTERVALL_STUNDEN"))
|
||||
if os.environ.get("SPIELE_NEUHEITEN_QUELLEN_SYNC_AKTIV", "1").strip() == "1":
|
||||
self._quellen_scheduler_starten(
|
||||
os.environ.get("SPIELE_NEUHEITEN_QUELLEN_SYNC_INTERVALL_STUNDEN")
|
||||
)
|
||||
|
||||
def on_unload(self) -> None:
|
||||
if self._scheduler is not None:
|
||||
@@ -248,7 +338,34 @@ class NeuheitenPlugin(BasePlugin):
|
||||
finally:
|
||||
self._sync_sperre.release()
|
||||
|
||||
def _scheduler_starten(self, intervall_stunden_roh: str | None) -> None:
|
||||
def _quellen_sync_ausfuehren(self) -> list[QuellenLaufZeile]:
|
||||
"""Ein Lauf über alle aktiven Web-Quellen (Fehler pro Quelle isoliert)."""
|
||||
dedup = (
|
||||
self.context.registry.get("dedup")
|
||||
if self.context is not None and self.context.registry is not None
|
||||
else None
|
||||
)
|
||||
service = WebQuellenSyncService(self.context.session_factory, dedup=dedup)
|
||||
return service.laufe()
|
||||
|
||||
def _quellen_sync_job(self) -> None:
|
||||
"""Hintergrund-Job: holt regelmäßig die Web-Quellen ab."""
|
||||
if not self._quellen_sperre.acquire(blocking=False):
|
||||
_logger.info("Web-Quellen-Sync läuft bereits — Durchlauf übersprungen.")
|
||||
return
|
||||
try:
|
||||
zeilen = self._quellen_sync_ausfuehren()
|
||||
text = " · ".join(z.als_text() for z in zeilen)
|
||||
if any(z.fehlermeldung for z in zeilen):
|
||||
_logger.warning("Web-Quellen-Sync mit Fehlern: %s", text)
|
||||
else:
|
||||
_logger.info("Web-Quellen-Sync abgeschlossen: %s", text)
|
||||
except Exception as exc:
|
||||
_logger.warning("Web-Quellen-Sync fehlgeschlagen: %s", exc)
|
||||
finally:
|
||||
self._quellen_sperre.release()
|
||||
|
||||
def _quellen_scheduler_starten(self, intervall_stunden_roh: str | None) -> None:
|
||||
from apscheduler.schedulers.background import BackgroundScheduler
|
||||
|
||||
try:
|
||||
@@ -257,7 +374,40 @@ class NeuheitenPlugin(BasePlugin):
|
||||
intervall_stunden = 24
|
||||
intervall_stunden = max(intervall_stunden, 1)
|
||||
|
||||
scheduler = BackgroundScheduler()
|
||||
scheduler = self._scheduler_oder_neu()
|
||||
erster_lauf = datetime.now(timezone.utc) + timedelta(seconds=15)
|
||||
scheduler.add_job(
|
||||
self._quellen_sync_job,
|
||||
trigger="interval",
|
||||
hours=intervall_stunden,
|
||||
next_run_time=erster_lauf,
|
||||
id="web-quellen-sync",
|
||||
replace_existing=True,
|
||||
)
|
||||
_logger.info(
|
||||
"Web-Quellen-Sync aktiv: alle %s h (%s)",
|
||||
intervall_stunden,
|
||||
", ".join(k.anzeigename for k in ADAPTER_KLASSEN),
|
||||
)
|
||||
|
||||
def _scheduler_oder_neu(self):
|
||||
"""Ein gemeinsamer BackgroundScheduler für alle Sync-Jobs des Plugins."""
|
||||
from apscheduler.schedulers.background import BackgroundScheduler
|
||||
|
||||
if self._scheduler is None or not self._scheduler.running:
|
||||
scheduler = BackgroundScheduler()
|
||||
scheduler.start()
|
||||
self._scheduler = scheduler
|
||||
return self._scheduler
|
||||
|
||||
def _scheduler_starten(self, intervall_stunden_roh: str | None) -> None:
|
||||
try:
|
||||
intervall_stunden = int(intervall_stunden_roh or "24")
|
||||
except ValueError:
|
||||
intervall_stunden = 24
|
||||
intervall_stunden = max(intervall_stunden, 1)
|
||||
|
||||
scheduler = self._scheduler_oder_neu()
|
||||
erster_lauf = datetime.now(timezone.utc) + timedelta(seconds=10)
|
||||
scheduler.add_job(
|
||||
self._sync_job,
|
||||
@@ -267,8 +417,6 @@ class NeuheitenPlugin(BasePlugin):
|
||||
id="bgg-neuheiten-sync",
|
||||
replace_existing=True,
|
||||
)
|
||||
scheduler.start()
|
||||
self._scheduler = scheduler
|
||||
_logger.info(
|
||||
"BGG-Neuheiten-Sync aktiv: alle %s h, Suchbegriffe: %s",
|
||||
intervall_stunden,
|
||||
|
||||
Reference in New Issue
Block a user