"""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, ) from .uploaders import notify_email, upload_folder, upload_nextcloud, upload_sftp log = logging.getLogger(__name__) class PreflightError(RuntimeError): """Erforderliche externe Binaries fehlen.""" # 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_preflight(pdfa_level: str = "", skip_text: bool = False) -> 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`). 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) 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 run(self) -> None: check_preflight(self.cfg.ocr.pdfa_level, self.cfg.ocr.skip_text) check_output_config(self.cfg.output.original_on_success, self.cfg.output.archive_dir, self.cfg.output.name_mode) 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: while not self._stop.is_set(): self._stop.wait(1.0) finally: self.shutdown() def run_once(self) -> int: """Verarbeitet alle bereits liegenden PDFs (incoming/ + working/) und beendet sich. Returns: Anzahl fehlgeschlagener PDFs (0 = alles ok). """ check_preflight(self.cfg.ocr.pdfa_level, self.cfg.ocr.skip_text) check_output_config(self.cfg.output.original_on_success, self.cfg.output.archive_dir, self.cfg.output.name_mode) 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() for p in sorted(self.cfg.paths.incoming.iterdir()): if _is_pdf(p): self.enqueue(p) 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: 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)