cd803a3dfe
Vor dem Rollout durchgesehen und die verbliebenen Stellen geschlossen, an denen etwas schiefgehen konnte, ohne dass es irgendwo sichtbar wurde. Datenverlust: - veraPDF: das in [verapdf].binary konfigurierte Programm wird im Preflight geprueft. Bisher galt bei falschem Pfad JEDE Datei als "nicht konform" — Ergebnis nach error/, Original geloescht (Default delete). run_verapdf() trennt jetzt ausserdem ein echtes FAIL-Urteil von einer Stoerung (VeraPdfUnavailable: nicht startbar, abgestuerzt, kein PASS/FAIL in der Ausgabe). Bei Stoerung wandern Original UND Ergebnis nach error/, das Original wird nicht entsorgt. - Gleichnamige Dateien wurden in outgoing/, error/ und beim Ordner-Upload mit abweichendem target kommentarlos ueberschrieben. Jetzt Zeitstempel daneben, mit Warnung; ProcessResult.output traegt den echten Pfad. Robustheit: - Kaputtes oder nicht lesbares TOML beim Start: Exit 2 statt Traceback. - RestartPreventExitStatus=2 in der Unit — Exit 2 (Config/Preflight) laeuft nicht mehr endlos neu, die Instanz bleibt sichtbar failed stehen. - Toter watchdog-Observer wird erkannt: Exit 3, systemd setzt den Watch neu auf. Vorher blieb die Unit "active" und verarbeitete nichts mehr. - Relative Pfade in [paths]/archive_dir/target sind ein Config-Fehler statt still unter /opt zu landen. - Fehler beim Archivieren entwertet den Durchlauf nicht mehr: Upload und Mail laufen, Sichtbarkeit ueber log.error + "OK mit Warnung"-Mail. - Nicht-PDFs in incoming/ werden beim Start-Scan gesammelt gemeldet. - Logging explizit nach stdout (die Doku versprach das schon). Struktur: - Neue lib/common.sh, von install.sh und update.sh gesourct. Die doppelte venv_is_healthy() gibt es nur noch einmal, in der gruendlichen Fassung — die schlanke in install.sh haette eine nach einem Distro-Sprung kaputte venv als gesund durchgewunken (nachgewiesen). - install.sh warnt in Containern, wenn systemd-journald nicht laeuft. Doku: Dateisystem-Festlegung (ext4/xfs/zfs, kein CIFS/NFS wegen inotify), Debian 13 in LXC auf Proxmox scheitert an journald (243/CREDENTIALS, AppArmor blockiert sd-mkdcreds) inkl. Abhilfe, echte Speicher-Messwerte, Exit-Code-Tabelle. 254 Tests gruen (vorher 152). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
692 lines
26 KiB
Python
692 lines
26 KiB
Python
"""Hauptservice: Hotfolder via watchdog, ThreadPool für PDF-Verarbeitung."""
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import re
|
||
import shutil
|
||
import signal
|
||
import subprocess
|
||
import threading
|
||
import time
|
||
from concurrent.futures import Future, ThreadPoolExecutor
|
||
from datetime import datetime
|
||
from pathlib import Path
|
||
|
||
from watchdog.events import FileSystemEvent, FileSystemEventHandler
|
||
from watchdog.observers import Observer
|
||
|
||
from .config import Config
|
||
from .processor import (
|
||
OCR_TEMP_PREFIX,
|
||
VALID_NAME_MODES,
|
||
ProcessResult,
|
||
_move_to_error,
|
||
process_pdf,
|
||
resolve_verapdf_binary,
|
||
)
|
||
from .uploaders import notify_email, upload_folder, upload_nextcloud, upload_sftp
|
||
|
||
log = logging.getLogger(__name__)
|
||
|
||
|
||
class PreflightError(RuntimeError):
|
||
"""Erforderliche externe Binaries fehlen."""
|
||
|
||
|
||
# Exit-Code, mit dem sich der Dienst bei totem watchdog-Observer beendet.
|
||
# Bewusst NICHT 2: die Unit setzt RestartPreventExitStatus=2 für Config- und
|
||
# Preflight-Fehler, die ein Neustart nicht heilt. Ein toter Observer soll
|
||
# dagegen genau das — neu starten, damit der inotify-Watch neu aufgesetzt wird.
|
||
EXIT_OBSERVER_DEAD = 3
|
||
|
||
|
||
# Pflicht-Binaries für ocrmypdf
|
||
_REQUIRED_BINARIES = ("tesseract", "gs")
|
||
|
||
# Ghostscript-Versionen mit bekanntem Bug (Issue #3):
|
||
# 10.0.0 .. 10.02.0 (inklusive). Ab 10.02.1 wieder nutzbar.
|
||
_GS_BROKEN_MIN = (10, 0, 0)
|
||
_GS_BROKEN_MAX = (10, 2, 0)
|
||
|
||
# Ab dieser ocrmypdf-Major steht die Ghostscript-Pruefung hinter
|
||
# `if options.output_type.startswith('pdfa')` — mit pdfa_level = "" wird
|
||
# Ghostscript gar nicht angefasst und die Pruefung greift nicht.
|
||
# Darunter (16.x und aelter) laeuft sie BEDINGUNGSLOS, also auch bei
|
||
# output_type="pdf": dort reicht skip_text=true, um auf Debian 12 jede
|
||
# einzelne PDF scheitern zu lassen. Quelle jeweils
|
||
# ocrmypdf/builtin_plugins/ghostscript.py::check_options().
|
||
_OCRMYPDF_GS_GUARD_MAJOR = 17
|
||
|
||
|
||
def _parse_version(text: str) -> tuple[int, ...] | None:
|
||
"""Extrahiert die erste X.Y[.Z] Version aus einem String."""
|
||
m = re.search(r"(\d+)\.(\d+)(?:\.(\d+))?", text)
|
||
if not m:
|
||
return None
|
||
return tuple(int(x) if x is not None else 0 for x in m.groups())
|
||
|
||
|
||
def is_ghostscript_broken(version: str | None) -> bool:
|
||
"""Prüft, ob eine Ghostscript-Version vom bekannten Bug betroffen ist.
|
||
|
||
Betrifft 10.0.0 bis einschließlich 10.02.0. Ab 10.02.1 wieder sicher.
|
||
"""
|
||
if not version:
|
||
return False
|
||
parsed = _parse_version(version)
|
||
if parsed is None:
|
||
return False
|
||
# Auf 3-Tupel normalisieren
|
||
while len(parsed) < 3:
|
||
parsed = parsed + (0,)
|
||
parsed = parsed[:3]
|
||
return _GS_BROKEN_MIN <= parsed <= _GS_BROKEN_MAX
|
||
|
||
|
||
def detect_ghostscript_version() -> str | None:
|
||
"""Ruft `gs --version` auf und gibt den Versionsstring zurück (oder None)."""
|
||
gs = shutil.which("gs")
|
||
if gs is None:
|
||
return None
|
||
try:
|
||
result = subprocess.run([gs, "--version"], capture_output=True,
|
||
text=True, timeout=5)
|
||
except (OSError, subprocess.TimeoutExpired):
|
||
return None
|
||
return result.stdout.strip() or None
|
||
|
||
|
||
def detect_ocrmypdf_version() -> str | None:
|
||
"""Liest die installierte ocrmypdf-Version aus den Paket-Metadaten.
|
||
|
||
Bewusst über `importlib.metadata` statt über einen Import: das ist
|
||
billiger und funktioniert auch in den Tests, in denen ocrmypdf gar nicht
|
||
installiert ist (dann None).
|
||
"""
|
||
try:
|
||
from importlib.metadata import version
|
||
return version("ocrmypdf")
|
||
except Exception: # noqa: BLE001 - fehlende Metadaten dürfen nichts umwerfen
|
||
return None
|
||
|
||
|
||
def ocrmypdf_checks_gs_always(version: str | None) -> bool:
|
||
"""True, wenn ocrmypdf die Ghostscript-Pruefung unabhaengig vom output_type fährt.
|
||
|
||
Das ist bei 16.x und aelter der Fall (siehe `_OCRMYPDF_GS_GUARD_MAJOR`).
|
||
Ist die Version unbekannt, wird `False` angenommen: der Pin in
|
||
requirements.txt steht auf 17.x, und ein Fehlalarm, der den Dienst nicht
|
||
starten laesst, waere schlimmer als die fehlende Warnung.
|
||
"""
|
||
if not version:
|
||
return False
|
||
parsed = _parse_version(version)
|
||
if parsed is None:
|
||
return False
|
||
return parsed[0] < _OCRMYPDF_GS_GUARD_MAJOR
|
||
|
||
|
||
def check_output_config(mode: str, archive_dir: str,
|
||
name_mode: str = "prefix") -> None:
|
||
"""Validiert die [output]-Section. Wirft PreflightError bei Problemen."""
|
||
valid_modes = {"delete", "archive"}
|
||
if mode not in valid_modes:
|
||
raise PreflightError(
|
||
f"[output].original_on_success={mode!r} ungültig. "
|
||
f"Erlaubt: {sorted(valid_modes)}"
|
||
)
|
||
if mode == "archive" and not archive_dir:
|
||
raise PreflightError(
|
||
"[output].original_on_success='archive' erfordert [output].archive_dir"
|
||
)
|
||
# Früh prüfen: sonst schlägt ein Tippfehler erst pro Datei zu — und zwar
|
||
# NACH dem Move nach working/, wo die Datei dann liegen bleibt.
|
||
if name_mode not in VALID_NAME_MODES:
|
||
raise PreflightError(
|
||
f"[output].name_mode={name_mode!r} ungültig. "
|
||
f"Erlaubt: {sorted(VALID_NAME_MODES)}"
|
||
)
|
||
|
||
|
||
def check_verapdf_binary(enabled: bool, binary: str) -> None:
|
||
"""Prüft das in [verapdf].binary konfigurierte Programm — wenn aktiviert.
|
||
|
||
Ohne diese Prüfung ist ein Tippfehler im Pfad der gefährlichste Fehler des
|
||
ganzen Dienstes: `run_verapdf()` findet das Programm für JEDE Datei nicht,
|
||
das OCR-Ergebnis wandert nach error/, und `_dispose_original()` löscht bei
|
||
`original_on_success = "delete"` (dem Default) das Original. Scan für Scan
|
||
verschwinden so die Vorlagen, während die Unit als `active (running)`
|
||
dasteht.
|
||
"""
|
||
if not enabled:
|
||
return
|
||
if not binary:
|
||
raise PreflightError(
|
||
"[verapdf].enabled = true, aber [verapdf].binary ist leer. "
|
||
"Entweder den Pfad zum veraPDF-Programm eintragen oder "
|
||
"[verapdf].enabled = false setzen."
|
||
)
|
||
if resolve_verapdf_binary(binary) is None:
|
||
raise PreflightError(
|
||
f"[verapdf].enabled = true, aber [verapdf].binary = {binary!r} "
|
||
"existiert nicht oder ist nicht ausführbar. Der Dienst startet "
|
||
"bewusst nicht: ein nicht aufrufbares veraPDF würde sonst jede "
|
||
"einzelne PDF als ungültig werten, das OCR-Ergebnis nach error/ "
|
||
"schieben und das Original laut [output].original_on_success "
|
||
"entsorgen. Pfad korrigieren (chmod +x nicht vergessen) oder "
|
||
"[verapdf].enabled = false setzen."
|
||
)
|
||
|
||
|
||
def check_preflight(pdfa_level: str = "", skip_text: bool = False,
|
||
verapdf_enabled: bool = False,
|
||
verapdf_binary: str = "") -> None:
|
||
"""Prüft externe Abhängigkeiten.
|
||
|
||
- Tesseract und Ghostscript müssen im PATH sein
|
||
- Die Ghostscript-Version wird gegen den bekannten 10.0.0–10.02.0 Bug
|
||
geprüft, und zwar genau unter der Bedingung, unter der ocrmypdf selbst
|
||
abbricht (siehe `_gs_block_reason`).
|
||
- Ist [verapdf].enabled gesetzt, muss auch das dort konfigurierte
|
||
Programm vorhanden und ausführbar sein (siehe `check_verapdf_binary`).
|
||
|
||
Wirft PreflightError bei fehlenden Binaries oder unsicherem Ghostscript.
|
||
"""
|
||
missing = [b for b in _REQUIRED_BINARIES if shutil.which(b) is None]
|
||
if missing:
|
||
raise PreflightError(
|
||
"Fehlende Abhängigkeiten: " + ", ".join(missing)
|
||
+ ". Bitte installieren: sudo apt install tesseract-ocr ghostscript"
|
||
)
|
||
|
||
reason = _gs_block_reason(pdfa_level, skip_text)
|
||
if reason:
|
||
raise PreflightError(reason)
|
||
|
||
check_verapdf_binary(verapdf_enabled, verapdf_binary)
|
||
|
||
|
||
def _gs_block_reason(pdfa_level: str, skip_text: bool) -> str | None:
|
||
"""Liefert die Fehlermeldung, wenn ocrmypdf mit diesem Ghostscript abbricht.
|
||
|
||
Abgebildet wird die reale Bedingung aus
|
||
`ocrmypdf/builtin_plugins/ghostscript.py::check_options()`:
|
||
|
||
betroffene GS-Version UND (skip_text ODER redo_ocr)
|
||
UND (PDF/A-Ausgabe ODER ocrmypdf < 17)
|
||
|
||
Die letzte Klammer ist der Teil, der v0.6.0 durchrutschen ließ: bis
|
||
einschließlich ocrmypdf 16.x steht die Pruefung ohne jeden Guard in
|
||
`check_options()` und schlaegt deshalb auch bei `output_type="pdf"` zu.
|
||
Ab 17.0.0 umschliesst sie ein
|
||
`if options.output_type.startswith('pdfa'):` — ohne PDF/A wird
|
||
Ghostscript nicht angefasst.
|
||
|
||
`redo_ocr` kennt unsere Config nicht (es gibt keinen entsprechenden Key in
|
||
`OcrConfig`), deshalb steht es hier bewusst nicht in der Bedingung.
|
||
|
||
Returns:
|
||
Fehlermeldung oder None, wenn die Kombination unkritisch ist.
|
||
"""
|
||
if not skip_text:
|
||
# Weder skip_text noch redo_ocr — ocrmypdf fasst den Pfad nicht an.
|
||
return None
|
||
|
||
gs_version = detect_ghostscript_version()
|
||
if not is_ghostscript_broken(gs_version):
|
||
return None
|
||
|
||
ocrmypdf_version = detect_ocrmypdf_version()
|
||
always = ocrmypdf_checks_gs_always(ocrmypdf_version)
|
||
if not pdfa_level and not always:
|
||
return None
|
||
|
||
if pdfa_level:
|
||
ursache = (
|
||
f"[ocr].pdfa_level = {pdfa_level!r} (PDF/A-Ausgabe) zusammen mit "
|
||
"[ocr].skip_text = true"
|
||
)
|
||
else:
|
||
ursache = (
|
||
f"[ocr].skip_text = true und ocrmypdf {ocrmypdf_version} — bis "
|
||
f"einschließlich {_OCRMYPDF_GS_GUARD_MAJOR - 1}.x prüft ocrmypdf "
|
||
"Ghostscript auch dann, wenn gar kein PDF/A erzeugt wird. Jede "
|
||
"einzelne PDF würde in error/ landen"
|
||
)
|
||
|
||
return (
|
||
f"Ghostscript {gs_version} ist von einem bekannten Fehler betroffen "
|
||
"(10.0.0–10.02.0, der Debian-12-Standard) und wird von ocrmypdf "
|
||
f"abgelehnt: {ursache}. "
|
||
"Abhilfe — eines von beidem: "
|
||
"(1) Ghostscript >= 10.02.1 aus bookworm-backports installieren "
|
||
"(install.sh bietet das an): "
|
||
"echo 'deb http://deb.debian.org/debian bookworm-backports main' | "
|
||
"sudo tee /etc/apt/sources.list.d/bookworm-backports.list && "
|
||
"sudo apt update && sudo apt install -t bookworm-backports ghostscript — "
|
||
"oder (2) in der Config [ocr].skip_text = false setzen "
|
||
"(dann wird vorhandener Text neu erkannt statt übersprungen)"
|
||
+ (" bzw. [ocr].pdfa_level = \"\"." if pdfa_level else ".")
|
||
)
|
||
|
||
|
||
def _is_pdf(path: Path) -> bool:
|
||
return path.suffix.lower() == ".pdf" and path.is_file()
|
||
|
||
|
||
def _wait_until_stable(path: Path, checks: int = 3, interval: float = 1.0) -> bool:
|
||
"""Wartet bis Datei nicht mehr wächst (Scanner schreibt mehrmals)."""
|
||
last = -1
|
||
stable_count = 0
|
||
for _ in range(60): # max ~60s
|
||
try:
|
||
size = path.stat().st_size
|
||
except FileNotFoundError:
|
||
return False
|
||
if size == last and size > 0:
|
||
stable_count += 1
|
||
if stable_count >= checks:
|
||
return True
|
||
else:
|
||
stable_count = 0
|
||
last = size
|
||
time.sleep(interval)
|
||
return False
|
||
|
||
|
||
class _Handler(FileSystemEventHandler):
|
||
def __init__(self, service: "HotfolderService") -> None:
|
||
self.service = service
|
||
|
||
def on_created(self, event: FileSystemEvent) -> None:
|
||
if not event.is_directory:
|
||
self.service.enqueue(Path(event.src_path))
|
||
|
||
def on_moved(self, event: FileSystemEvent) -> None:
|
||
if not event.is_directory:
|
||
self.service.enqueue(Path(event.dest_path))
|
||
|
||
def on_closed(self, event: FileSystemEvent) -> None:
|
||
if not event.is_directory:
|
||
self.service.enqueue(Path(event.src_path))
|
||
|
||
|
||
class HotfolderService:
|
||
def __init__(self, cfg: Config) -> None:
|
||
self.cfg = cfg
|
||
self._executor = ThreadPoolExecutor(
|
||
max_workers=cfg.ocr.max_workers,
|
||
thread_name_prefix="ocr",
|
||
)
|
||
self._observer: Observer | None = None
|
||
self._stop = threading.Event()
|
||
self._inflight: set[str] = set()
|
||
self._lock = threading.Lock()
|
||
self._success_count = 0
|
||
self._error_count = 0
|
||
|
||
@property
|
||
def success_count(self) -> int:
|
||
return self._success_count
|
||
|
||
@property
|
||
def error_count(self) -> int:
|
||
return self._error_count
|
||
|
||
# ---- Setup ----
|
||
|
||
def ensure_dirs(self) -> None:
|
||
for p in (self.cfg.paths.incoming, self.cfg.paths.outgoing,
|
||
self.cfg.paths.working, self.cfg.paths.error):
|
||
p.mkdir(parents=True, exist_ok=True)
|
||
|
||
# ---- Lifecycle ----
|
||
|
||
def _preflight(self) -> None:
|
||
check_preflight(self.cfg.ocr.pdfa_level, self.cfg.ocr.skip_text,
|
||
self.cfg.verapdf.enabled, self.cfg.verapdf.binary)
|
||
check_output_config(self.cfg.output.original_on_success,
|
||
self.cfg.output.archive_dir,
|
||
self.cfg.output.name_mode)
|
||
|
||
def run(self) -> int:
|
||
"""Startet den Dienst und läuft, bis gestoppt wird.
|
||
|
||
Returns:
|
||
0 bei regulärem Stopp (SIGTERM/SIGINT), sonst `EXIT_OBSERVER_DEAD`.
|
||
"""
|
||
self._preflight()
|
||
self.ensure_dirs()
|
||
self._scan_existing()
|
||
|
||
self._observer = Observer()
|
||
self._observer.schedule(_Handler(self), str(self.cfg.paths.incoming), recursive=False)
|
||
self._observer.start()
|
||
log.info("Hotfolder läuft. Watching: %s", self.cfg.paths.incoming)
|
||
|
||
signal.signal(signal.SIGTERM, lambda *_: self._stop.set())
|
||
signal.signal(signal.SIGINT, lambda *_: self._stop.set())
|
||
|
||
try:
|
||
return self._wait_loop()
|
||
finally:
|
||
self.shutdown()
|
||
|
||
def _wait_loop(self) -> int:
|
||
"""Hauptschleife: wartet auf den Stopp und bewacht den Observer.
|
||
|
||
Stirbt der watchdog-Observer im Betrieb (erschöpftes
|
||
inotify-Watch-Limit, ersetztes oder neu gemountetes Verzeichnis),
|
||
blieb die Unit bisher `active (running)` und verarbeitete nichts mehr:
|
||
kein Log, keine Mail, niemand merkt es. Für einen Hotfolder ist das
|
||
der schlechteste denkbare Zustand. Deshalb wird der Observer
|
||
sekündlich mitgeprüft und der Dienst im Ernstfall mit
|
||
`EXIT_OBSERVER_DEAD` beendet, damit systemd ihn per
|
||
`Restart=on-failure` neu startet und den Watch neu aufsetzt.
|
||
"""
|
||
while not self._stop.is_set():
|
||
self._stop.wait(1.0)
|
||
if self._stop.is_set():
|
||
# Regulärer Stopp — hier darf kein Fehlalarm entstehen, auch
|
||
# wenn der Observer planmäßig schon gestoppt wurde.
|
||
break
|
||
if self._observer is not None and not self._observer.is_alive():
|
||
log.error(
|
||
"Der Verzeichnis-Watch auf %s ist gestorben — es werden "
|
||
"KEINE neuen Dateien mehr erkannt. Mögliche Ursachen: "
|
||
"erschöpftes inotify-Watch-Limit "
|
||
"(fs.inotify.max_user_watches), ersetztes oder neu "
|
||
"gemountetes Verzeichnis. Der Dienst beendet sich mit "
|
||
"Exit %d, damit systemd ihn neu startet und der Watch "
|
||
"neu aufgesetzt wird.",
|
||
self.cfg.paths.incoming, EXIT_OBSERVER_DEAD,
|
||
)
|
||
return EXIT_OBSERVER_DEAD
|
||
return 0
|
||
|
||
def run_once(self) -> int:
|
||
"""Verarbeitet alle bereits liegenden PDFs (incoming/ + working/) und beendet sich.
|
||
|
||
Returns:
|
||
Anzahl fehlgeschlagener PDFs (0 = alles ok).
|
||
"""
|
||
self._preflight()
|
||
self.ensure_dirs()
|
||
self._scan_existing()
|
||
self._executor.shutdown(wait=True)
|
||
log.info("One-shot fertig: %d ok, %d Fehler",
|
||
self._success_count, self._error_count)
|
||
return self._error_count
|
||
|
||
def shutdown(self) -> None:
|
||
log.info("Shutdown läuft...")
|
||
if self._observer:
|
||
self._observer.stop()
|
||
self._observer.join(timeout=5)
|
||
self._executor.shutdown(wait=True, cancel_futures=False)
|
||
log.info("Shutdown ok.")
|
||
|
||
# ---- Queue ----
|
||
|
||
def _scan_existing(self) -> None:
|
||
"""Beim Start: bereits liegende PDFs aufgreifen.
|
||
|
||
Zuerst working/ (abgebrochene Läufe, siehe `_scan_working`), danach
|
||
incoming/. Die Reihenfolge ist wichtig, damit eine Namenskollision
|
||
zwischen beiden Verzeichnissen aufgelöst ist, bevor die
|
||
incoming-Datei nach working/ will.
|
||
"""
|
||
self._scan_working()
|
||
fremd: list[str] = []
|
||
for p in sorted(self.cfg.paths.incoming.iterdir()):
|
||
if _is_pdf(p):
|
||
self.enqueue(p)
|
||
elif p.is_file():
|
||
fremd.append(p.name)
|
||
self._report_non_pdf(fremd)
|
||
|
||
def _report_non_pdf(self, names: list[str]) -> None:
|
||
"""Meldet einmalig, wie viele Fremddateien in incoming/ liegen.
|
||
|
||
Alles ohne .pdf-Endung wird ignoriert und sammelte sich bisher stumm
|
||
an — Scanner-Fehlablagen, abgebrochene Uploads, Thumbnails. Eine
|
||
Sammelmeldung beim Start-Scan, keine Zeile pro Datei und nichts im
|
||
laufenden Betrieb: das soll auffallen, nicht spammen.
|
||
"""
|
||
if not names:
|
||
return
|
||
beispiele = ", ".join(names[:3])
|
||
if len(names) > 3:
|
||
beispiele += f", … (+{len(names) - 3} weitere)"
|
||
log.warning(
|
||
"In %s liegen %d Datei(en) ohne .pdf-Endung — sie werden nicht "
|
||
"verarbeitet und bleiben dort liegen: %s",
|
||
self.cfg.paths.incoming, len(names), beispiele,
|
||
)
|
||
|
||
def _scan_working(self) -> None:
|
||
"""Greift Dateien auf, die ein harter Stopp in working/ liegen ließ.
|
||
|
||
`process_pdf()` verschiebt das Original vor dem OCR nach working/.
|
||
Wird der Dienst dort abgeschossen (SIGKILL nach TimeoutStopSec),
|
||
bleibt es liegen und wurde bisher nie wieder angefasst — stiller
|
||
Datenverlust. Die Datei wird deshalb an Ort und Stelle
|
||
wiederaufgenommen; `process_pdf()` erkennt das und verschiebt sie
|
||
nicht erneut.
|
||
|
||
Die Zwischendateien des abgebrochenen OCR-Laufs (Präfix `__ocr_`)
|
||
sind unvollständige Fragmente: als Eingabe unbrauchbar und als
|
||
Ergebnis wertlos. Sie werden gelöscht, damit sie niemand für ein
|
||
fertiges PDF hält und damit der neue Lauf sauber startet.
|
||
"""
|
||
working = self.cfg.paths.working
|
||
if not working.is_dir():
|
||
return
|
||
for p in sorted(working.iterdir()):
|
||
if not p.is_file():
|
||
continue
|
||
if p.name.startswith(OCR_TEMP_PREFIX):
|
||
log.warning(
|
||
"Unvollständiges OCR-Fragment aus abgebrochenem Lauf "
|
||
"gefunden und gelöscht: %s", p,
|
||
)
|
||
try:
|
||
p.unlink()
|
||
except OSError:
|
||
log.exception("Konnte OCR-Fragment %s nicht löschen", p)
|
||
continue
|
||
if not _is_pdf(p):
|
||
continue
|
||
target = self._free_resume_name(p)
|
||
log.warning(
|
||
"Abgebrochener Lauf wird fortgesetzt: %s lag noch in %s "
|
||
"(Dienst wurde vermutlich hart gestoppt) — OCR startet neu",
|
||
target.name, working,
|
||
)
|
||
self.enqueue(target)
|
||
|
||
def _free_resume_name(self, p: Path) -> Path:
|
||
"""Entschärft eine Namenskollision zwischen working/ und incoming/.
|
||
|
||
Liegt in incoming/ eine gleichnamige (aber andere) Datei, würden beide
|
||
dieselbe working- und dieselbe outgoing-Datei beanspruchen. Die
|
||
wiederaufgenommene Datei bekommt deshalb einen Zeitstempel angehängt —
|
||
dann laufen beide durch, statt dass eine überschrieben wird.
|
||
"""
|
||
if not (self.cfg.paths.incoming / p.name).exists():
|
||
return p
|
||
ts = datetime.now().strftime("%Y%m%d-%H%M%S")
|
||
renamed = p.with_name(f"{p.stem}_{ts}{p.suffix}")
|
||
try:
|
||
p.rename(renamed)
|
||
except OSError:
|
||
log.exception("Konnte %s nicht umbenennen — Wiederaufnahme unter "
|
||
"Originalnamen", p)
|
||
return p
|
||
log.warning(
|
||
"In %s liegt eine gleichnamige Datei %s — die wiederaufgenommene "
|
||
"Datei wurde nach %s umbenannt, damit sich beide nicht "
|
||
"überschreiben", self.cfg.paths.incoming, p.name, renamed.name,
|
||
)
|
||
return renamed
|
||
|
||
def enqueue(self, path: Path) -> None:
|
||
if not _is_pdf(path):
|
||
return
|
||
key = str(path.resolve())
|
||
with self._lock:
|
||
if key in self._inflight:
|
||
return
|
||
self._inflight.add(key)
|
||
fut = self._executor.submit(self._process, path)
|
||
fut.add_done_callback(lambda f, k=key: self._done(k, f))
|
||
|
||
def _done(self, key: str, fut: Future) -> None:
|
||
with self._lock:
|
||
self._inflight.discard(key)
|
||
exc = fut.exception()
|
||
if exc:
|
||
log.exception("Worker-Exception", exc_info=exc)
|
||
|
||
# ---- Processing ----
|
||
|
||
def _count_success(self) -> None:
|
||
with self._lock:
|
||
self._success_count += 1
|
||
|
||
def _count_error(self) -> None:
|
||
with self._lock:
|
||
self._error_count += 1
|
||
|
||
def _process(self, path: Path) -> None:
|
||
if not _wait_until_stable(path):
|
||
if not path.exists():
|
||
# Datei wurde währenddessen entfernt — kein Fehlerfall
|
||
log.info("Datei vor der Verarbeitung verschwunden: %s", path)
|
||
return
|
||
# Bewusst als Fehler zählen: sonst liefert --once trotz liegen
|
||
# gebliebener Datei Exit 0.
|
||
log.error(
|
||
"Datei hat sich nicht stabilisiert (Timeout): %s — bleibt in %s "
|
||
"liegen und wird beim nächsten Lauf erneut versucht",
|
||
path, self.cfg.paths.incoming,
|
||
)
|
||
self._count_error()
|
||
return
|
||
if not path.exists():
|
||
return
|
||
|
||
try:
|
||
result: ProcessResult = process_pdf(
|
||
src=path,
|
||
working_dir=self.cfg.paths.working,
|
||
outgoing_dir=self.cfg.paths.outgoing,
|
||
error_dir=self.cfg.paths.error,
|
||
ocr_cfg=self.cfg.ocr,
|
||
vera_cfg=self.cfg.verapdf,
|
||
output_cfg=self.cfg.output,
|
||
)
|
||
except Exception as e: # noqa: BLE001 - kein Fehler darf die Zählung umgehen
|
||
log.exception("Unerwarteter Fehler bei der Verarbeitung von %s", path.name)
|
||
self._count_error()
|
||
self._rescue_to_error(path)
|
||
self._notify(ProcessResult(
|
||
path, self.cfg.paths.outgoing / path.name, False,
|
||
f"unerwarteter Fehler: {e}",
|
||
))
|
||
return
|
||
|
||
if not result.success:
|
||
self._count_error()
|
||
self._notify(result)
|
||
return
|
||
|
||
failed = self._dispatch_uploads(result.output)
|
||
if failed:
|
||
log.error(
|
||
"Upload fehlgeschlagen (%s) für %s — das OCR selbst war "
|
||
"erfolgreich, die Datei bleibt daher in %s liegen und wird "
|
||
"NICHT nach error/ verschoben",
|
||
", ".join(failed), result.output.name, result.output.parent,
|
||
)
|
||
self._count_error()
|
||
self._notify_upload_failure(result, failed)
|
||
return
|
||
|
||
self._count_success()
|
||
self._notify(result)
|
||
|
||
def _rescue_to_error(self, src: Path) -> None:
|
||
"""Bringt eine Datei nach einer unerwarteten Exception ins error-Verzeichnis.
|
||
|
||
Die Datei kann je nach Abbruchzeitpunkt noch in incoming/ oder schon in
|
||
working/ liegen. Der erste Treffer wird verschoben (keine Doppel-Moves),
|
||
Fehler beim Verschieben werden nur geloggt.
|
||
"""
|
||
error_dir = self.cfg.paths.error
|
||
for candidate in (src, self.cfg.paths.working / src.name):
|
||
try:
|
||
if not candidate.is_file():
|
||
continue
|
||
if candidate.parent.resolve() == error_dir.resolve():
|
||
return # liegt bereits im error-Verzeichnis
|
||
except OSError:
|
||
continue
|
||
_move_to_error(candidate, error_dir)
|
||
return
|
||
log.warning("Datei %s nach Fehler nicht mehr auffindbar — "
|
||
"kein Verschieben nach error/ möglich", src.name)
|
||
|
||
def _dispatch_uploads(self, pdf: Path) -> list[str]:
|
||
"""Schiebt das fertige PDF an alle Upload-Ziele.
|
||
|
||
Die uploader prüfen `cfg.enabled` jeweils selbst und liefern für
|
||
deaktivierte Ziele True.
|
||
|
||
Returns:
|
||
Namen der fehlgeschlagenen Ziele — leere Liste = alle erfolgreich.
|
||
"""
|
||
failed: list[str] = []
|
||
if not upload_folder(pdf, self.cfg.folder, self.cfg.paths.outgoing):
|
||
failed.append("folder")
|
||
if not upload_nextcloud(pdf, self.cfg.nextcloud):
|
||
failed.append("nextcloud")
|
||
if not upload_sftp(pdf, self.cfg.sftp):
|
||
failed.append("sftp")
|
||
return failed
|
||
|
||
def _notify_upload_failure(self, result: ProcessResult, failed: list[str]) -> None:
|
||
"""Fehler-Mail, wenn das OCR lief, aber mindestens ein Upload scheiterte."""
|
||
subject = f"[pdf-ocr] FEHLER Upload: {result.source.name}"
|
||
body = (
|
||
f"OCR erfolgreich: {result.output}\n\n"
|
||
f"Fehlgeschlagene Upload-Ziele: {', '.join(failed)}\n\n"
|
||
f"Das OCR-PDF bleibt in {result.output.parent} liegen und wurde "
|
||
"NICHT nach error/ verschoben. Details siehe Log.\n"
|
||
)
|
||
notify_email(self.cfg.email, subject, body, False)
|
||
|
||
def _notify(self, result: ProcessResult) -> None:
|
||
if result.success and result.warning:
|
||
# Erfolgreich verarbeitet, aber das Original blieb liegen. Der
|
||
# Durchlauf zählt als Erfolg (das PDF ist fertig und ausgeliefert),
|
||
# die Mail geht aber als Nicht-Erfolg raus, damit sie auch bei
|
||
# [notify.email].on = "errors" zugestellt wird — sonst wäre das
|
||
# genau wieder ein stiller Fehlerpfad.
|
||
subject = f"[pdf-ocr] OK mit Warnung: {result.source.name}"
|
||
body = (
|
||
f"Datei verarbeitet: {result.output}\n\n"
|
||
f"ACHTUNG: {result.warning}\n"
|
||
)
|
||
notify_email(self.cfg.email, subject, body, False)
|
||
return
|
||
if result.success:
|
||
subject = f"[pdf-ocr] OK: {result.source.name}"
|
||
body = f"Datei verarbeitet: {result.output}\n"
|
||
if result.verapdf_passed is not None:
|
||
body += f"veraPDF: {'PASS' if result.verapdf_passed else 'FAIL'}\n"
|
||
else:
|
||
subject = f"[pdf-ocr] FEHLER: {result.source.name}"
|
||
body = f"Fehler beim Verarbeiten von {result.source}\n\n{result.error}\n"
|
||
notify_email(self.cfg.email, subject, body, result.success)
|