542 lines
23 KiB
Python
542 lines
23 KiB
Python
"""Plugin „neuheiten“ — Neuheitenliste aus der BoardGameGeek XML API2.
|
|
|
|
Umfang:
|
|
- Eigene Migration (Tabelle `neuheiten`).
|
|
- BGG-Client mit Rate-Limit/Retry (siehe bgg.py), Sync-Service (sync.py).
|
|
- Regelmäßiger Sync als Hintergrund-Job (APScheduler, konfigurierbar).
|
|
- Manueller „Jetzt synchronisieren“-Button für Admins/Redakteure.
|
|
- Sortier-/filterbare Neuheitenliste mit Suche (deutsche UI).
|
|
|
|
Konfiguration über Umgebungsvariablen (gelesen beim Plugin-Start;
|
|
SPIELE_BGG_TOKEN wird dagegen vor jedem Sync gelesen):
|
|
- SPIELE_BGG_SYNC_AKTIV „1“ (Standard) = Hintergrund-Job an
|
|
- SPIELE_BGG_SYNC_INTERVALL_STUNDEN Intervall in Stunden (Standard 24)
|
|
- SPIELE_BGG_SUCHBEGRIFFE Komma-getrennte Suchbegriffe
|
|
(Standard: „brettspiel“)
|
|
- SPIELE_BGG_MAX_TREFFER_PRO_SUCHE Obergrenze Treffer/Suchbegriff (25)
|
|
- SPIELE_BGG_TOKEN API-Token für die BGG-XML-API2
|
|
(Pflicht seit BGG-Umstellung, wird als
|
|
„Authorization: Bearer …“ gesendet);
|
|
ohne Token wird der Sync übersprungen
|
|
|
|
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
|
|
|
|
from fastapi import Depends, Form, HTTPException, Request
|
|
from fastapi.responses import RedirectResponse
|
|
from sqlalchemy import func, or_, select
|
|
|
|
from redaktionskern.auth.deps import require_roles, require_user
|
|
from redaktionskern.auth.models import Role, User
|
|
from redaktionskern.contracts import BasePlugin, Migration, NavEntry, PluginContext
|
|
|
|
from .bgg import BggClient, BggFehler
|
|
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")
|
|
|
|
MELDUNG_OHNE_TOKEN = (
|
|
"Kein BGG-API-Token konfiguriert (SPIELE_BGG_TOKEN) — Sync übersprungen."
|
|
)
|
|
|
|
SORTIERBAR = {
|
|
"titel": Neuheit.titel,
|
|
"verlag": Neuheit.verlag,
|
|
"autor": Neuheit.autor,
|
|
"jahr": Neuheit.erscheinungsjahr,
|
|
"status": Neuheit.status,
|
|
"aktualisiert": Neuheit.aktualisiert_am,
|
|
}
|
|
STANDARD_SORTIERUNG = ("titel", "auf")
|
|
MAX_MELDUNGS_LAENGE = 300
|
|
|
|
|
|
def _umgebung_liste(name: str, standard: str) -> list[str]:
|
|
import os
|
|
|
|
roh = os.environ.get(name, standard)
|
|
return [teil.strip() for teil in roh.split(",") if teil.strip()]
|
|
|
|
|
|
class NeuheitenPlugin(BasePlugin):
|
|
name = "neuheiten"
|
|
title = "Neuheiten"
|
|
description = (
|
|
"Neuheitenliste aus der BoardGameGeek-API — regelmäßiger Sync als "
|
|
"Hintergrund-Job, Erweiterungen und Prototypen werden gefiltert."
|
|
)
|
|
|
|
def __init__(self) -> None:
|
|
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
|
|
|
|
@self.router.get("/neuheiten")
|
|
def seite(
|
|
request: Request,
|
|
user: User = Depends(require_user),
|
|
q: str = "",
|
|
status: str = "",
|
|
sort: str = STANDARD_SORTIERUNG[0],
|
|
richtung: str = STANDARD_SORTIERUNG[1],
|
|
meldung: str = "",
|
|
fehler: str = "",
|
|
seite_nr: int = 1,
|
|
):
|
|
spalte = SORTIERBAR.get(sort, SORTIERBAR[STANDARD_SORTIERUNG[0]])
|
|
abwaerts = richtung == "ab" if sort in SORTIERBAR else False
|
|
spalte_sortiert = spalte.desc() if abwaerts else spalte.asc()
|
|
pro_seite = 50
|
|
|
|
with self.context.session_factory() as db:
|
|
abfrage = select(Neuheit)
|
|
if q.strip():
|
|
muster = f"%{q.strip().lower()}%"
|
|
abfrage = abfrage.where(
|
|
or_(
|
|
func.lower(Neuheit.titel).like(muster),
|
|
func.lower(func.coalesce(Neuheit.verlag, "")).like(muster),
|
|
func.lower(func.coalesce(Neuheit.autor, "")).like(muster),
|
|
)
|
|
)
|
|
if status.strip():
|
|
abfrage = abfrage.where(Neuheit.status == status.strip())
|
|
gesamt = db.scalar(select(func.count()).select_from(abfrage.subquery())) or 0
|
|
eintraege = db.scalars(
|
|
abfrage.order_by(spalte_sortiert, Neuheit.id.asc())
|
|
.offset((max(seite_nr, 1) - 1) * pro_seite)
|
|
.limit(pro_seite)
|
|
).all()
|
|
status_optionen = [
|
|
zeile for zeile in db.scalars(
|
|
select(Neuheit.status).distinct().order_by(Neuheit.status)
|
|
)
|
|
]
|
|
|
|
quellen_zeilen = self._quellen_uebersicht(db)
|
|
hat_weitere = max(seite_nr, 1) * pro_seite < gesamt
|
|
kontext = {
|
|
"user": user,
|
|
"titel": self.title,
|
|
"eintraege": eintraege,
|
|
"q": q,
|
|
"status_filter": status,
|
|
"sort": sort if sort in SORTIERBAR else STANDARD_SORTIERUNG[0],
|
|
"richtung": "ab" if abwaerts else "auf",
|
|
"status_optionen": status_optionen,
|
|
"meldung": meldung[:MAX_MELDUNGS_LAENGE],
|
|
"fehler": fehler[:MAX_MELDUNGS_LAENGE],
|
|
"ist_redaktion": user.role in (
|
|
Role.ADMIN.value,
|
|
Role.REDAKTEUR.value,
|
|
),
|
|
"suchbegriffe": ", ".join(self._suchbegriffe) or "—",
|
|
"sync_aktiv": self._scheduler is not None and self._scheduler.running,
|
|
"anzahl": gesamt,
|
|
"hat_weitere": hat_weitere,
|
|
"naechste_seite": max(seite_nr, 1) + 1,
|
|
"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(
|
|
request=request,
|
|
name="neuheiten/_liste.html",
|
|
context=kontext,
|
|
)
|
|
return self.context.templates.TemplateResponse(
|
|
request=request,
|
|
name="neuheiten/index.html",
|
|
context=kontext,
|
|
)
|
|
|
|
@self.router.post("/neuheiten/sync")
|
|
def jetzt_synchronisieren(
|
|
user: User = Depends(require_roles(Role.ADMIN.value, Role.REDAKTEUR.value)),
|
|
suchbegriff: str = Form(""),
|
|
):
|
|
begriffe = [suchbegriff] if suchbegriff.strip() else self._suchbegriffe
|
|
try:
|
|
ergebnis = self._sync_ausfuehren(begriffe)
|
|
except BggFehler as exc:
|
|
_logger.warning("Manueller BGG-Sync fehlgeschlagen: %s", exc)
|
|
ziel = f"/neuheiten?meldung={quote(f'Sync fehlgeschlagen: {exc}')}"
|
|
return RedirectResponse(ziel, status_code=303)
|
|
if ergebnis.abbruch:
|
|
# z. B. fehlender Token oder HTTP 401 — klar benennbare Ursache
|
|
ziel = f"/neuheiten?fehler={quote(ergebnis.abbruch)}"
|
|
else:
|
|
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)),
|
|
quelle: str = Form(""),
|
|
):
|
|
"""Gesamt-Sync oder Einzel-Quellen-Sync (`quelle` = Quellname)."""
|
|
if quelle:
|
|
zeilen = self._quellen_sync_ausfuehren(nur_quellen={quelle})
|
|
ziel = f"/neuheiten?meldung={quote('Web-Quelle „' + quelle + '“ synchronisiert: ' + ' · '.join(z.als_text() for z in zeilen))}"
|
|
else:
|
|
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)
|
|
|
|
@self.router.post("/neuheiten/quellen/{quellen_name}/schalter")
|
|
def quellen_schalter(
|
|
quellen_name: str,
|
|
user: User = Depends(require_roles(Role.ADMIN.value, Role.REDAKTEUR.value)),
|
|
aktiv: str = Form(...),
|
|
):
|
|
"""An/Aus-Schalter einer Web-Quelle (übersteuert den Env-Standard)."""
|
|
if quellen_name not in {k.name for k in ADAPTER_KLASSEN}:
|
|
raise HTTPException(status_code=404, detail="Unbekannte Quelle")
|
|
assert self.context is not None and self.context.session_factory is not None
|
|
with self.context.session_factory() as db:
|
|
eintrag = db.scalar(
|
|
select(QuellenStatus).where(QuellenStatus.quelle == quellen_name)
|
|
)
|
|
if eintrag is None:
|
|
eintrag = QuellenStatus(quelle=quellen_name, url="")
|
|
db.add(eintrag)
|
|
eintrag.aktiv_schalter = aktiv == "1"
|
|
db.commit()
|
|
zustand = "eingeschaltet" if aktiv == "1" else "ausgeschaltet"
|
|
ziel = f"/neuheiten?meldung={quote(f'Quelle „{quellen_name}“ {zustand}.')}"
|
|
return RedirectResponse(ziel, 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)
|
|
schalter = status.aktiv_schalter if status is not None else None
|
|
aktiv = (
|
|
schalter
|
|
if schalter is not None
|
|
else quelle_aktiv(klasse.name)
|
|
)
|
|
ansichten.append(
|
|
{
|
|
"name": klasse.name,
|
|
"anzeigename": klasse.anzeigename,
|
|
"url": klasse.start_urls[0] if klasse.start_urls else "",
|
|
"aktiv": aktiv,
|
|
"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)
|
|
|
|
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)
|
|
|
|
def bild_url_spalte(conn) -> None:
|
|
"""Migration 0004: Spalte bild_url für Coverbilder (nullable).
|
|
|
|
`ALTER TABLE … ADD COLUMN` mit nullable VARCHAR ohne Default ist
|
|
SQLite- und Postgres-portabel und erhält bestehende Zeilen.
|
|
"""
|
|
info = conn.exec_driver_sql("PRAGMA table_info(neuheiten)").fetchall()
|
|
vorhanden = any(zeile[1] == "bild_url" for zeile in info)
|
|
if conn.dialect.name == "sqlite" and vorhanden:
|
|
return
|
|
if conn.dialect.name != "sqlite":
|
|
existiert = conn.exec_driver_sql(
|
|
"SELECT column_name FROM information_schema.columns "
|
|
"WHERE table_name = 'neuheiten' AND column_name = 'bild_url'"
|
|
).scalar()
|
|
if existiert is not None:
|
|
return
|
|
conn.exec_driver_sql(
|
|
"ALTER TABLE neuheiten ADD COLUMN bild_url VARCHAR(500)"
|
|
)
|
|
|
|
def quellen_url_spalte(conn) -> None:
|
|
"""Migration 0005: Spalte quellen_url für den Link zur Quelle."""
|
|
info = conn.exec_driver_sql("PRAGMA table_info(neuheiten)").fetchall()
|
|
vorhanden = any(zeile[1] == "quellen_url" for zeile in info)
|
|
if conn.dialect.name == "sqlite" and vorhanden:
|
|
return
|
|
if conn.dialect.name != "sqlite":
|
|
existiert = conn.exec_driver_sql(
|
|
"SELECT column_name FROM information_schema.columns "
|
|
"WHERE table_name = 'neuheiten' AND column_name = 'quellen_url'"
|
|
).scalar()
|
|
if existiert is not None:
|
|
return
|
|
conn.exec_driver_sql(
|
|
"ALTER TABLE neuheiten ADD COLUMN quellen_url VARCHAR(500)"
|
|
)
|
|
|
|
def aktiv_schalter_spalte(conn) -> None:
|
|
"""Migration 0006: UI-Schalter pro Web-Quelle (nullable Boolean)."""
|
|
info = conn.exec_driver_sql(
|
|
"PRAGMA table_info(neuheiten_quellen_status)"
|
|
).fetchall()
|
|
vorhanden = any(zeile[1] == "aktiv_schalter" for zeile in info)
|
|
if conn.dialect.name == "sqlite" and vorhanden:
|
|
return
|
|
if conn.dialect.name != "sqlite":
|
|
existiert = conn.exec_driver_sql(
|
|
"SELECT column_name FROM information_schema.columns "
|
|
"WHERE table_name = 'neuheiten_quellen_status' "
|
|
"AND column_name = 'aktiv_schalter'"
|
|
).scalar()
|
|
if existiert is not None:
|
|
return
|
|
conn.exec_driver_sql(
|
|
"ALTER TABLE neuheiten_quellen_status ADD COLUMN aktiv_schalter BOOLEAN"
|
|
)
|
|
|
|
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),
|
|
Migration(version="0004_bild_url", up=bild_url_spalte),
|
|
Migration(version="0005_quellen_url", up=quellen_url_spalte),
|
|
Migration(version="0006_aktiv_schalter", up=aktiv_schalter_spalte),
|
|
]
|
|
|
|
def navigation(self) -> list[NavEntry]:
|
|
return [NavEntry(label=self.title, url="/neuheiten")]
|
|
|
|
def on_load(self, context: PluginContext) -> None:
|
|
super().on_load(context)
|
|
import os
|
|
|
|
self._suchbegriffe = _umgebung_liste("SPIELE_BGG_SUCHBEGRIFFE", "brettspiel")
|
|
try:
|
|
self._max_treffer_pro_suche = int(
|
|
os.environ.get("SPIELE_BGG_MAX_TREFFER_PRO_SUCHE", "25")
|
|
)
|
|
except ValueError:
|
|
self._max_treffer_pro_suche = 25
|
|
|
|
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:
|
|
self._scheduler.shutdown(wait=False)
|
|
self._scheduler = None
|
|
super().on_unload()
|
|
|
|
# ---------- Internas ----------
|
|
|
|
def _aktives_bgg_token(self) -> str | None:
|
|
"""Liest den BGG-API-Token aus der Umgebung (leer = nicht gesetzt).
|
|
|
|
Bewusst pro Sync gelesen, damit ein rotierter Token ohne Neustart
|
|
greift.
|
|
"""
|
|
import os
|
|
|
|
return (os.environ.get("SPIELE_BGG_TOKEN") or "").strip() or None
|
|
|
|
def _neuer_client(self) -> BggClient:
|
|
"""Fabrik für den BGG-Client; von Tests überschreibbar."""
|
|
return BggClient(token=self._aktives_bgg_token())
|
|
|
|
def _sync_ausfuehren(self, suchbegriffe: list[str]):
|
|
if self._aktives_bgg_token() is None:
|
|
# Ohne Token lehnt die BGG-API jeden Request mit 401 ab —
|
|
# den Sync gar nicht erst starten.
|
|
_logger.warning(MELDUNG_OHNE_TOKEN)
|
|
return SyncErgebnis(abbruch=MELDUNG_OHNE_TOKEN)
|
|
client = self._neuer_client()
|
|
try:
|
|
service = SyncService(
|
|
self.context.session_factory,
|
|
client,
|
|
max_treffer_pro_suchbegriff=self._max_treffer_pro_suche,
|
|
)
|
|
return service.synchronisiere(suchbegriffe)
|
|
finally:
|
|
client.schliessen()
|
|
|
|
def _sync_job(self) -> None:
|
|
"""Hintergrund-Job: holt regelmäßig die Neuheiten ab."""
|
|
if not self._sync_sperre.acquire(blocking=False):
|
|
_logger.info("BGG-Sync läuft bereits — Durchlauf übersprungen.")
|
|
return
|
|
try:
|
|
ergebnis = self._sync_ausfuehren(self._suchbegriffe)
|
|
if ergebnis.abbruch:
|
|
_logger.warning("BGG-Sync nicht ausgeführt: %s", ergebnis.abbruch)
|
|
else:
|
|
_logger.info("BGG-Sync abgeschlossen: %s", ergebnis.als_text())
|
|
except Exception as exc:
|
|
_logger.warning("BGG-Sync fehlgeschlagen: %s", exc)
|
|
finally:
|
|
self._sync_sperre.release()
|
|
|
|
def _quellen_sync_ausfuehren(
|
|
self, nur_quellen: set[str] | None = None
|
|
) -> list[QuellenLaufZeile]:
|
|
"""Ein Lauf über alle aktiven Web-Quellen (Fehler pro Quelle isoliert).
|
|
|
|
`nur_quellen`: optional nur diese Quellnamen (manueller Einzellauf,
|
|
ignoriert den An/Aus-Schalter der Quelle).
|
|
"""
|
|
registry = (
|
|
self.context.registry
|
|
if self.context is not None and self.context.registry is not None
|
|
else None
|
|
)
|
|
dedup = registry.get("dedup") if registry is not None else None
|
|
audit = registry.get("audit-log") if registry is not None else None
|
|
service = WebQuellenSyncService(
|
|
self.context.session_factory, dedup=dedup, audit=audit
|
|
)
|
|
return service.laufe(nur_quellen=nur_quellen)
|
|
|
|
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:
|
|
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=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,
|
|
trigger="interval",
|
|
hours=intervall_stunden,
|
|
next_run_time=erster_lauf,
|
|
id="bgg-neuheiten-sync",
|
|
replace_existing=True,
|
|
)
|
|
_logger.info(
|
|
"BGG-Neuheiten-Sync aktiv: alle %s h, Suchbegriffe: %s",
|
|
intervall_stunden,
|
|
", ".join(self._suchbegriffe),
|
|
)
|
|
|
|
|
|
plugin = NeuheitenPlugin()
|