diff --git a/.github/workflows/docker-worker.yml b/.github/workflows/docker-worker.yml index 4b57efd..786ae7a 100644 --- a/.github/workflows/docker-worker.yml +++ b/.github/workflows/docker-worker.yml @@ -1,12 +1,24 @@ +# Purpose: Build and publish reviewed worker images for changed runtime sources. +# Input/Output: Reads source files and writes immutable multi-architecture images to GHCR. +# Invariants: Documentation-only and Compose-pin commits do not create a new image digest. +# Debugging: Inspect the Build and push worker image step plus its published digest. + name: Docker Worker Publish on: workflow_dispatch: push: branches: - - main + - "**" tags: - "v*" + paths: + - "src/**" + - "requirements.txt" + - "docker/Dockerfile" + - "docker/entrypoint.sh" + - "custom_components/paperless_kiplus/manifest.json" + - ".github/workflows/docker-worker.yml" permissions: contents: read @@ -19,6 +31,14 @@ jobs: - name: Checkout uses: actions/checkout@v4 + - name: Read application version + id: app + shell: bash + run: | + set -euo pipefail + VERSION=$(python3 -c 'import json; print(json.load(open("custom_components/paperless_kiplus/manifest.json", encoding="utf-8"))["version"])') + echo "version=$VERSION" >> "$GITHUB_OUTPUT" + - name: Set up QEMU uses: docker/setup-qemu-action@v3 @@ -52,3 +72,6 @@ jobs: push: true tags: ${{ steps.meta.outputs.tags }} labels: ${{ steps.meta.outputs.labels }} + build-args: | + APP_COMMIT=${{ github.sha }} + APP_VERSION=${{ steps.app.outputs.version }} diff --git a/.github/workflows/python-tests.yml b/.github/workflows/python-tests.yml new file mode 100644 index 0000000..3aec6c4 --- /dev/null +++ b/.github/workflows/python-tests.yml @@ -0,0 +1,48 @@ +# Purpose: Run the local Python, manifest, Compose, and image checks in CI. +# Input/Output: Reads the repository at the pushed commit and publishes only test results. +# Invariants: No deployment, production secret, or remote Paperless access is required. +# Debugging: Re-run the first failed command locally from the repository root. + +name: Python + Docker Tests + +on: + pull_request: + push: + +permissions: + contents: read + +jobs: + test: + name: Python 3.12 + runs-on: ubuntu-latest + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Set up Python + uses: actions/setup-python@v5 + with: + python-version: "3.12" + cache: pip + + - name: Install test dependencies + run: | + python -m pip install --upgrade pip + python -m pip install -r requirements.txt pytest + + - name: Compile Python sources + run: >- + python -m py_compile + src/*.py + custom_components/paperless_kiplus/*.py + tests/*.py + + - name: Run unit and integration tests + run: python -m pytest -q + + - name: Validate production Compose + run: docker compose -f docker/docker-compose.unraid-broker.yml config --quiet + + - name: Build production worker image + run: docker build --file docker/Dockerfile --tag paperless-kiplus-worker:test . diff --git a/CHANGELOG.md b/CHANGELOG.md new file mode 100644 index 0000000..d401f7c --- /dev/null +++ b/CHANGELOG.md @@ -0,0 +1,26 @@ +# Changelog + +## Unreleased + +### Added + +- Persistente, idempotente HTTP-202-Hintergrundjobs für Sorter, Restart, + Entity-Scan und Entity-Merge. +- Authentifizierter Jobstatus mit sicheren Request-IDs, exaktem Fortschritt + und Reload-Recovery in beiden Weboberflächen. +- Restart-, Parallelitäts-, Redaction-, API- und Docker-CI-Tests. + +### Changed + +- Home Assistant pollt Remote-Jobs bis zum terminalen Status. +- Browser-Token werden nur noch sitzungsbezogen gespeichert. +- Lange Live-Review-Abfragen verwenden den Job-Endpunkt. + +### Security + +- API-, Log- und Konfigurationsantworten sind nicht cachebar. +- Jobfehler und Worker-Logs maskieren Zugangsdaten und Providerdetails. +- Unterbrochene Schreibjobs werden nicht automatisch wiederholt. +- Python-Basisimage und Runtime-Abhängigkeiten sind reproduzierbar gepinnt. +- Produktion baut commitgebunden im Broker und benötigt keinen privaten + Registry-Pull. diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 62c3496..a6db9a6 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -17,12 +17,30 @@ Vielen Dank für dein Interesse an Beiträgen zu **Paperless KIplus**. ## Lokale Checks -Vor einem PR bitte mindestens: - -1. Syntax prüfen: - - `python3 -m py_compile custom_components/paperless_kiplus/*.py src/paperless_ai_sorter.py` -2. Integration laden und einen Testlauf in Home Assistant durchführen -3. Prüfen, dass README/Docs bei neuen Features aktualisiert sind +Vor einem PR bitte aus dem Repository-Root ausführen: + +```bash +python3 -m pip install -r requirements.txt pytest +python3 -m pytest -q +docker compose -f docker/docker-compose.unraid-broker.yml config --quiet +docker build -f docker/Dockerfile -t paperless-kiplus-worker:test . +``` + +Bei Änderungen an Hintergrundjobs zusätzlich mindestens einen Happy Path, +einen ungültigen Input, Restart-Recovery und fünf wiederholte parallele +Admission-Läufe testen. Echte Paperless- oder API-Token dürfen nie in Fixtures, +Logs oder Fehlermeldungen erscheinen. + +Für gezieltes Debugging: + +```bash +PAPERLESS_KIPLUS_LOG_LEVEL=DEBUG python3 src/worker_api.py --data-dir ./worker-data +python3 -m unittest tests.test_background_jobs -v +``` + +Anschließend die Integration in Home Assistant laden und einen kleinen Dry-Run +gegen eine Testinstanz durchführen. README und Betriebsdoku müssen das neue +Verhalten erklären. ## Pull-Request Ablauf diff --git a/README.md b/README.md index 3db7619..a47441d 100644 --- a/README.md +++ b/README.md @@ -901,13 +901,15 @@ Was der Worker mitbringt: - vollständige Ausführung ohne Home Assistant - eingebaute Weboberfläche unter `/` - Entity-Review-Seite unter `/review` für Dopplungen und KI-Lernhinweise -- JSON-API für Run / Stop / Resume / Restart / Backfill +- persistente HTTP-202-Jobs für Run / Resume / Restart / Backfill und Review +- Idempotency-Keys, sichere Request-IDs und Reload-Recovery mit Backoff - Log-Download, Status und Konfigurationsverwaltung - persistente Dateien für Config, Metriken und Resume-State unter `/data` Dokumentation: - [Docker- und Unraid-Betrieb](./docs/docker-unraid.md) +- [Cloudflare-sichere Langläufer](./docs/cloudflare-long-running-jobs.md) - [Migration von Home Assistant zum Remote-Worker](./docs/migration-ha-to-worker.md) - [Lokale LLMs für kleinere Aufgaben](./docs/local-llm-routing.md) @@ -925,76 +927,38 @@ Danach: - Dopplungsreview: `http://:8787/review` - Status-API: `http://:8787/api/status` -### Robuste Unraid-Installation +### Sichere Unraid-Installation in der Feberdin-Umgebung -Für Unraid gibt es jetzt zwei klare Wege: +Die produktive Compose-Quelle liegt unter +`docker/docker-compose.unraid-broker.yml`. Deployments erfolgen ausschließlich +über den Unraid Deployment Broker und einen vollständigen Git-Commit-SHA: -1. Von macOS/Linux per SSH auf einen entfernten Unraid-Server deployen -2. Direkt im Unraid-Terminal ohne Repo-Checkout installieren -3. Direkt auf dem Unraid-Server mit vorhandenem Repo installieren +1. `stack_source_status` +2. `stack_validate` +3. `deploy_plan` +4. `approval_request`, falls der Plan dies verlangt +5. `deploy_apply` +6. `deployment_status`, `docker_list` und `logs_tail` -#### Empfohlen: Remote-Deploy von macOS/Linux nach Unraid +Der Stack erhält `PAPERLESS_KIPLUS_TOKEN` zur Laufzeit über +`secret://PAPERLESS_KIPLUS_TOKEN`. Echte Tokenwerte gehören weder in Git noch +in Chat, Logs oder Compose. Das bestehende Appdata-Verzeichnis +`/mnt/user/appdata/paperless-kiplus` wird unverändert als `/data` eingebunden. -```bash -bash docker/deploy-to-unraid.sh \ - --unraid-host 192.168.178.30 \ - --paperless-url http://192.168.178.20:8000 \ - --paperless-token PAPERLESS_TOKEN \ - --ai-api-key OPENAI_KEY \ - --ai-model gpt-4.1-mini -``` - -Das Remote-Skript: - -- verbindet sich per SSH mit Unraid -- kopiert den eigentlichen Host-Installer auf den Server -- überträgt optional eine lokale `config.yaml` -- führt die Installation direkt auf Unraid aus - -#### Direkte Ausführung auf dem Unraid-Server - -Wenn du direkt im Unraid-Terminal bist und das Repo dort nicht lokal liegen -hast, kannst du jetzt den Bootstrap-Weg nutzen. Er legt zuerst den passenden -Ordner an, lädt den Installer herunter und startet ihn direkt: - -```bash -mkdir -p /boot/config/custom/paperless-kiplus && \ -curl -fsSL https://raw.githubusercontent.com/Feberdin/Paperless-KIplus/v1.4.6/docker/bootstrap-unraid-worker.sh \ - -o /boot/config/custom/paperless-kiplus/bootstrap-unraid-worker.sh && \ -chmod +x /boot/config/custom/paperless-kiplus/bootstrap-unraid-worker.sh && \ -bash /boot/config/custom/paperless-kiplus/bootstrap-unraid-worker.sh \ - --ref v1.4.6 \ - --paperless-url http://192.168.178.20:8000 \ - --paperless-token PAPERLESS_TOKEN \ - --ai-api-key OPENAI_KEY \ - --ai-model gpt-4.1-mini -``` - -Dabei wird der Installer standardmäßig hier abgelegt: - -```text -/boot/config/custom/paperless-kiplus/install-unraid-worker.sh -``` - -Wenn du bereits ein Repo-Checkout auf Unraid hast, kannst du weiterhin direkt -das Host-Skript verwenden: +Der Broker baut das Worker-Image aus dem commitgebundenen Checkout. Dockerfile, +Basisimage, Python-Abhängigkeiten, App-Version und geprüfter App-Commit sind +festgelegt; damit ist kein privater Registry-Pull für die Produktion nötig. -```bash -bash /pfad/zum/repo/docker/install-unraid-worker.sh \ - --paperless-url http://192.168.178.20:8000 \ - --paperless-token PAPERLESS_TOKEN \ - --ai-api-key OPENAI_KEY \ - --ai-model gpt-4.1-mini -``` +Ein Rollback verwendet denselben Broker-Ablauf mit dem vorherigen GitOps-Commit. +Direkte SSH-, Docker-CLI- oder Unraid-Shell-Deployments sind für diese +Produktionsumgebung nicht vorgesehen. -Das Host-Skript: +### Logging und Fehlersuche -- legt die Appdata-Verzeichnisse an -- sichert bestehende Dateien -- erzeugt eine startfähige `config.yaml` -- schreibt einen Compose-Stack mit GHCR-Image -- startet oder aktualisiert den Worker -- prüft die API per Health-Check +Der Standard-Level ist `INFO`. Für einen zeitlich begrenzten Diagnose-Lauf kann +im Broker-Stack `PAPERLESS_KIPLUS_LOG_LEVEL=DEBUG` gesetzt werden. Logs werden +über `logs_tail` abgerufen; Zugangsdaten werden vor Datei-, Speicher- und +UI-Ausgabe maskiert. Jobfehler lassen sich über ihre `request_id` zuordnen. ### Remote-Steuerung aus Home Assistant diff --git a/custom_components/paperless_kiplus/manifest.json b/custom_components/paperless_kiplus/manifest.json index 1bbd699..1a78497 100644 --- a/custom_components/paperless_kiplus/manifest.json +++ b/custom_components/paperless_kiplus/manifest.json @@ -9,5 +9,5 @@ "iot_class": "local_polling", "issue_tracker": "https://github.com/Feberdin/Paperless-KIplus/issues", "requirements": [], - "version": "1.4.20" + "version": "1.4.21" } diff --git a/custom_components/paperless_kiplus/remote_runner.py b/custom_components/paperless_kiplus/remote_runner.py index bd3fbaa..2627ff5 100644 --- a/custom_components/paperless_kiplus/remote_runner.py +++ b/custom_components/paperless_kiplus/remote_runner.py @@ -31,10 +31,11 @@ from __future__ import annotations import asyncio -from dataclasses import dataclass -from datetime import UTC, datetime, timedelta import json import logging +import uuid +from dataclasses import dataclass +from datetime import UTC, datetime from pathlib import Path from typing import Any @@ -140,6 +141,8 @@ def __init__( self.last_config_sync_status: str = "idle" self._poll_task: asyncio.Task | None = None + self._active_job_status_url: str = "" + self._active_job_request_id: str = "" self._lock = asyncio.Lock() self._session = async_get_clientsession(hass) @@ -259,17 +262,21 @@ async def _api_json( path: str, *, payload: dict[str, Any] | None = None, + idempotency_key: str | None = None, ) -> dict[str, Any]: """Executes a JSON API request against the remote worker.""" if not self.remote_worker_url: raise ValueError("remote_worker_url ist nicht konfiguriert.") + headers = self._headers() + if idempotency_key: + headers["Idempotency-Key"] = idempotency_key timeout = aiohttp.ClientTimeout(total=60) async with self._session.request( method, self._worker_url(path), - headers=self._headers(), + headers=headers, json=payload, ssl=self.remote_worker_verify_ssl, timeout=timeout, @@ -284,7 +291,7 @@ async def _api_json( return {} parsed = json.loads(text) if not isinstance(parsed, dict): - raise RuntimeError(f"Remote-Worker API Antwort für {path} ist kein JSON-Objekt.") + raise TypeError(f"Remote-Worker API Antwort für {path} ist kein JSON-Objekt.") return parsed async def _api_text(self, path: str) -> str: @@ -373,11 +380,33 @@ def _apply_status_payload(self, payload: dict[str, Any]) -> None: self.last_log_export_path = str(payload.get("last_log_export_path") or self.last_log_export_path or "") worker_export_url = str(payload.get("last_log_export_url") or "").strip() if worker_export_url: - if worker_export_url.startswith("http://") or worker_export_url.startswith("https://"): + if worker_export_url.startswith(("http://", "https://")): self.last_log_export_url = worker_export_url else: self.last_log_export_url = self._worker_url(worker_export_url) + def _apply_action_response(self, payload: dict[str, Any]) -> None: + """Accept both legacy immediate responses and new HTTP-202 job admissions.""" + + status_url = str(payload.get("status_url") or "").strip() + if payload.get("job_id") and status_url: + self._active_job_status_url = status_url + self._active_job_request_id = str(payload.get("request_id") or "") + self.last_status = str(payload.get("status") or "queued") + self.last_message = ( + "Remote-Job angenommen" + + ( + f" (Request-ID: {self._active_job_request_id})" + if self._active_job_request_id + else "" + ) + ) + # The sorter process may need a short scheduling window. Keeping + # ``running`` true ensures HA polling does not stop before it starts. + self.running = True + return + self._apply_status_payload(payload.get("status") or payload) + async def _refresh_status(self) -> None: payload = await self._api_json("GET", "/api/status") self._apply_status_payload(payload) @@ -388,8 +417,33 @@ async def _poll_loop(self) -> None: try: while True: + active_job = bool(self._active_job_status_url) + terminal_failure: tuple[str, str] | None = None + job_status = "" + if active_job: + job = await self._api_json("GET", self._active_job_status_url) + job_status = str(job.get("status") or "") + if job_status in {"failed", "interrupted", "cancelled"}: + error = job.get("error") or {} + terminal_failure = ( + f"remote_job_{job_status}", + str(error.get("message") or "Remote-Hintergrundjob fehlgeschlagen."), + ) + self._active_job_status_url = "" + self._active_job_request_id = "" + elif job_status == "succeeded": + self._active_job_status_url = "" + self._active_job_request_id = "" await self._refresh_status() - if not self.running and not self.resume_available and self.last_status not in { + if terminal_failure: + self.running = False + self.last_status, self.last_message = terminal_failure + self.last_stderr_tail = self.last_message + self._notify() + if active_job and self._active_job_status_url and not self.running: + self.running = True + self.last_status = job_status or "queued" + if not self._active_job_status_url and not self.running and not self.resume_available and self.last_status not in { "waiting_auto_resume", "paused", }: @@ -398,9 +452,18 @@ async def _poll_loop(self) -> None: except asyncio.CancelledError: raise except Exception as exc: # noqa: BLE001 + request_hint = ( + f" Request-ID: {self._active_job_request_id}." + if self._active_job_request_id + else "" + ) self.last_status = "remote_poll_error" - self.last_message = str(exc) - self.last_stderr_tail = str(exc) + self.last_message = ( + "Remote-Jobstatus konnte nicht geladen werden; Worker-Verbindung prüfen." + + request_hint + ) + self.last_stderr_tail = self.last_message + _LOGGER.error("Remote-Jobstatus fehlgeschlagen: %s", type(exc).__name__) self._notify() def _ensure_polling(self) -> None: @@ -527,8 +590,9 @@ async def async_run( "max_documents": self.default_max_documents if max_documents is None else int(max_documents), "backfill_existing_documents": bool(backfill_existing_documents), }, + idempotency_key=uuid.uuid4().hex, ) - self._apply_status_payload(payload.get("status") or payload) + self._apply_action_response(payload) if self.running or self.last_status in {"waiting_auto_resume", "paused"}: self._ensure_polling() self._notify() @@ -549,8 +613,13 @@ async def async_force_stop(self) -> RunResult: return RunResult(self.last_status, self.last_exit_code, self.last_message) async def async_resume(self, *, force: bool = True) -> RunResult: - payload = await self._api_json("POST", "/api/resume", payload={"force": force}) - self._apply_status_payload(payload.get("status") or payload) + payload = await self._api_json( + "POST", + "/api/resume", + payload={"force": force}, + idempotency_key=uuid.uuid4().hex, + ) + self._apply_action_response(payload) self._ensure_polling() self._notify() return RunResult(self.last_status, self.last_exit_code, self.last_message) @@ -570,8 +639,9 @@ async def async_restart( "force": force, "backfill_existing_documents": backfill_existing_documents, }, + idempotency_key=uuid.uuid4().hex, ) - self._apply_status_payload(payload.get("status") or payload) + self._apply_action_response(payload) self._ensure_polling() self._notify() return RunResult(self.last_status, self.last_exit_code, self.last_message) diff --git a/docker/Dockerfile b/docker/Dockerfile index 5d35709..fa423ad 100644 --- a/docker/Dockerfile +++ b/docker/Dockerfile @@ -16,16 +16,24 @@ # - Open /api/status and /api/logs on the worker port. # - A startup error naming /data means the mounted directory is not writable by the runtime UID/GID. -FROM python:3.12.13-slim-trixie +# Why the digest is pinned: production may build directly from a broker-bound +# Git commit. Pinning the multi-architecture base manifest prevents a mutable +# registry tag from changing the resulting runtime between CI and Unraid. +FROM python:3.12.13-slim-trixie@sha256:229a2c5bfa27522db7815ea81f9bed70af17ccb9de9fc7ad142b1877b5830d36 ARG APP_UID=10001 ARG APP_GID=10001 +ARG APP_COMMIT=unknown +ARG APP_VERSION=dev ENV PYTHONDONTWRITEBYTECODE=1 \ PYTHONUNBUFFERED=1 \ PAPERLESS_KIPLUS_DATA_DIR=/data \ PAPERLESS_KIPLUS_HOST=0.0.0.0 \ - PAPERLESS_KIPLUS_PORT=8787 + PAPERLESS_KIPLUS_PORT=8787 \ + PAPERLESS_KIPLUS_LOG_LEVEL=INFO \ + PAPERLESS_KIPLUS_APP_VERSION="${APP_VERSION}" \ + PAPERLESS_KIPLUS_IMAGE_COMMIT="${APP_COMMIT}" WORKDIR /app diff --git a/docker/docker-compose.unraid-broker.yml b/docker/docker-compose.unraid-broker.yml index 53e2e5a..8cc3ece 100644 --- a/docker/docker-compose.unraid-broker.yml +++ b/docker/docker-compose.unraid-broker.yml @@ -12,8 +12,8 @@ # Important invariants: # - Do not replace the /data mount with a fresh relative directory; it contains # worker config, state, logs, metrics and exports. -# - Keep the worker image pinned to the reviewed multi-architecture manifest -# digest; update it only after the matching GitHub workflow succeeds. +# - Build only from this broker-bound Git checkout. APP_COMMIT names the exact +# source commit whose Docker build, tests, HACS, and Hassfest checks passed. # - Secrets stay as secret:// references and are injected only by the broker at # deploy_apply time. # - Heimdall labels expose only the dedicated safe metadata endpoint. @@ -29,7 +29,15 @@ x-broker: services: paperless-kiplus-worker: - image: ghcr.io/feberdin/paperless-kiplus-worker@sha256:04fc422c3025235ef5bd70742f2f9932521a32bfbfc487f2e0f45d5b7777eb96 + image: paperless-kiplus-worker:1.4.21-e155328 + build: + # The broker executes its generated Compose file from the checked-out + # stack root, so this path intentionally addresses the repository root. + context: . + dockerfile: docker/Dockerfile + args: + APP_COMMIT: e15532882f58a02e2a033d0f79e6b965eaad39a7 + APP_VERSION: 1.4.21 # Unraid appdata uses nobody:users by default. Keep the worker non-root # while preserving write access to the existing /data directory. user: "99:100" @@ -44,6 +52,7 @@ services: PAPERLESS_KIPLUS_DATA_DIR: /data PAPERLESS_KIPLUS_HOST: 0.0.0.0 PAPERLESS_KIPLUS_PORT: 8788 + PAPERLESS_KIPLUS_LOG_LEVEL: INFO PAPERLESS_KIPLUS_TOKEN: secret://PAPERLESS_KIPLUS_TOKEN volumes: - /mnt/user/appdata/paperless-kiplus:/data diff --git a/docs/cloudflare-long-running-jobs.md b/docs/cloudflare-long-running-jobs.md new file mode 100644 index 0000000..15a41a9 --- /dev/null +++ b/docs/cloudflare-long-running-jobs.md @@ -0,0 +1,164 @@ +# Cloudflare-sichere Langläufer + +## Zweck + +Dieses Dokument beschreibt den persistenten Job-Vertrag des Standalone-Workers, +die betroffenen Endpunkte und die sichere Migration. Der HTTP-Request nimmt nur +Arbeit an; Paperless-Zugriffe und Sorter-Prozesse laufen danach unabhängig vom +Cloudflare-Request weiter. + +Wichtige Invarianten: + +- Die HTTP-Annahme liefert schnell `202 Accepted`. +- Job- und Request-ID bleiben über Seiten-Reloads und Worker-Restarts erhalten. +- Es gibt keine erfundenen Prozentwerte oder Restzeiten. +- Schreibende Jobs laufen nie parallel auf derselben Ressource. +- Unterbrochene Schreibjobs werden nicht automatisch wiederholt. +- Parameter, Dokument-IDs, Rohfehler und Secrets erscheinen nicht im Jobstatus. +- `/api/status` und Heimdall melden Image-Version und vollständigen Build-Commit. + +Debugging erfolgt über `request_id`, den authentifizierten Jobstatus und die +redigierten Worker-Logs. `PAPERLESS_KIPLUS_LOG_LEVEL=DEBUG` aktiviert zusätzliche +Diagnosemeldungen. + +## Endpunkt-Inventar + +| Endpunkt | Bewertung | Verhalten | +|---|---|---| +| `POST /api/run` | potenziell lang | `202`, persistenter `sorter_run` | +| `POST /api/resume` | potenziell lang | `202`, persistenter `sorter_resume` | +| `POST /api/restart` | potenziell lang | `202`, ersetzt einen aktiven Sorter kontrolliert | +| `POST /api/review/entities/jobs` | potenziell lang, nur lesend | `202`, persistenter `review_scan` | +| `POST /api/review/merge` | potenziell lang, schreibend | `202`, persistenter `review_merge` | +| `GET /api/jobs/` | kurz | authentifizierter persistenter Status | +| `GET /api/review/entities` | früher blockierend | jetzt `405` mit Verweis auf den Job-Endpunkt | +| `GET /api/status`, `/api/logs`, `/api/config/*`, `/api/review/rules` | kurz | direkte lokale Antwort | +| `POST /api/stop`, `/api/stop_now` | kurz | setzt Stop-Signal bzw. terminiert den Prozess | +| `POST /api/config/import`, `/api/review/rules`, Reset-Endpunkte | kurz | validierter lokaler Dateizugriff | +| `/`, `/review`, `/api/heimdall/v1` | kurz | UI bzw. redigierter Health-Status | + +Der JSON-Body ist auf 1 MiB begrenzt. API-, Log- und Config-Antworten senden +`Cache-Control: no-store`, sodass Cloudflare und Browser private Daten nicht als +Cacheobjekt behandeln. + +## HTTP-Vertrag + +Ein Client kann optional einen maximal 200 Zeichen langen `Idempotency-Key` +senden. Dieselbe Operation mit demselben Schlüssel erhält erneut denselben Job. + +```http +POST /api/run +Authorization: Bearer +Idempotency-Key: +Content-Type: application/json + +{"dry_run":true,"max_documents":10} +``` + +```json +{ + "ok": true, + "status": "queued", + "job_id": "job_", + "request_id": "req_", + "status_url": "/api/jobs/job_", + "deduplicated": false +} +``` + +Der Status liefert `queued`, `running`, `succeeded`, `failed`, `interrupted` +oder `cancelled`. Unbekannte Werte bleiben explizit `null`: + +```json +{ + "status": "running", + "phase": "discovering_documents", + "progress_percent": null, + "estimated_seconds_remaining": null, + "attempt_count": 1, + "max_attempts": 1 +} +``` + +Terminale Fehler enthalten nur Fehlercode, sichere Handlungsanweisung und +`request_id`. Rohantworten externer APIs werden weder persistiert noch an die UI +zurückgegeben. + +## Persistenz, Parallelität und Restart + +Die additive SQLite-Datenbank liegt unter +`/data/state/background_jobs.sqlite3`. Transaktionen und ein partieller Unique +Index vergeben pro Ressourcenklasse genau einen aktiven Besitzer: + +- `paperless_write`: Sorter-Läufe und echte Entity-Merges +- `review_scan`: nur lesende Entity-Scans + +Ein Restart darf ausschließlich aktive Sorter-Jobs ersetzen. Ein gleichzeitig +laufender Merge führt weiterhin zu `409 Conflict`. Der Executor ist auf drei +Threads begrenzt; Jobausführungen haben genau einen Versuch. Bestehende, +ebenfalls begrenzte Paperless-HTTP-Retries bleiben davon unberührt. + +Beim Worker-Start gilt: + +- `review_scan` wird mit den minimal persistierten Parametern sicher neu geplant. +- Sorter- und Merge-Jobs werden `interrupted` und niemals automatisch erneut + ausgeführt. +- Ein unterbrochener Merge kann bereits einzelne Dokumente aktualisiert haben. + Vor einem manuellen Retry muss der Paperless-Zustand geprüft werden. +- Der Sorter behält zusätzlich seinen bestehenden Resume-State unter `/data/state`. + +Terminale Jobs werden nach 24 Stunden gelöscht; zusätzlich bleiben höchstens +200 terminale Datensätze erhalten. + +## Browser und Home Assistant + +Beide eingebetteten UIs speichern nur Job-, Request- und Status-URL im +`localStorage`. Das Bearer-Token liegt ausschließlich im `sessionStorage` und +wird nach dem Schließen des Tabs verworfen. Beim Reload werden aktive Jobs mit +exponentiellem Backoff von 750 ms bis maximal 10 s weiter beobachtet. + +Die Home-Assistant-Remote-Steuerung akzeptiert sowohl alte Sofortantworten als +auch den neuen 202-Vertrag. Neue Aktionen senden einen Idempotency-Key und +pollen die konkrete `status_url`, bis der Job terminal ist. + +## Cloudflare-Betrieb + +Ein separater Cloudflare Worker oder eine Durable Object Instanz ist nicht +erforderlich: Cloudflare sieht nur kurze Admission- und Status-Requests. Für +den Proxy müssen lediglich dieselben Pfade zum Origin durchgereicht und +Caching für API-Antworten deaktiviert bleiben; der Origin setzt dafür bereits +`no-store`. + +## Migration und Rollback + +Migration: + +1. Neues Image mit unverändertem `/data`-Mount deployen. +2. Worker-Health und `GET /api/status` prüfen. +3. Einen kleinen Dry-Run starten und die gelieferte `status_url` bis + `succeeded` pollen. +4. Browser neu laden und prüfen, dass der Job weiterhin angezeigt wird. + +Die SQLite-Tabelle wird beim Start additiv angelegt. Bestehende Config-, +Metrik-, Log- und Resume-Dateien werden nicht verschoben. + +Rollback: + +1. Im Broker den vorherigen GitOps-Commit planen; er enthält den zugehörigen + commitgebundenen Image-Build. +2. Plan prüfen und anwenden. +3. Die neue SQLite-Datei darf liegen bleiben; ältere Worker-Versionen ignorieren + sie. Für einen späteren Roll-forward bleibt dadurch die Diagnose erhalten. + +## Lokale Qualitätssicherung + +```bash +python3 -m pip install -r requirements.txt pytest +python3 -m pytest -q +docker compose -f docker/docker-compose.unraid-broker.yml config --quiet +docker build -f docker/Dockerfile -t paperless-kiplus-worker:test . +``` + +Negativfälle und fünfmal wiederholte Parallelitätsläufe sind Teil der +automatisierten Tests. Die produktive Bereitstellung erfolgt ausschließlich +über den Unraid Deployment Broker. diff --git a/docs/docker-unraid.md b/docs/docker-unraid.md index f6fb4bb..acb761e 100644 --- a/docs/docker-unraid.md +++ b/docs/docker-unraid.md @@ -25,111 +25,36 @@ Danach: - Web UI: `http://:8787/` - API-Status: `http://:8787/api/status` -## Schnellstart auf Unraid +## Produktion auf Unraid -### Einfachster Weg: Direkt im Unraid-Terminal ohne Repo-Checkout +In der Feberdin-Umgebung wird ausschließlich die GitOps-Quelle +`docker/docker-compose.unraid-broker.yml` über den Unraid Deployment Broker +bereitgestellt. Direkte SSH-, Shell-, Docker-CLI- und HTTP-Deployments sind +nicht Teil dieses Betriebswegs. -Wenn du direkt auf dem Unraid-Server arbeitest, kannst du den Installer jetzt -ohne lokales Git-Checkout starten. Das Bootstrap-Skript legt zuerst den -passenden Ordner an, laedt den eigentlichen Installer von GitHub und fuehrt ihn -anschliessend lokal auf Unraid aus: +Voraussetzungen: -```bash -mkdir -p /boot/config/custom/paperless-kiplus && \ -curl -fsSL https://raw.githubusercontent.com/Feberdin/Paperless-KIplus/v1.4.6/docker/bootstrap-unraid-worker.sh \ - -o /boot/config/custom/paperless-kiplus/bootstrap-unraid-worker.sh && \ -chmod +x /boot/config/custom/paperless-kiplus/bootstrap-unraid-worker.sh && \ -bash /boot/config/custom/paperless-kiplus/bootstrap-unraid-worker.sh \ - --ref v1.4.6 \ - --paperless-url http://192.168.178.20:8000 \ - --paperless-token PAPERLESS_TOKEN \ - --ai-api-key OPENAI_KEY \ - --ai-model gpt-4.1-mini -``` - -Der Installer landet standardmaessig hier: - -```text -/boot/config/custom/paperless-kiplus/install-unraid-worker.sh -``` - -### Empfohlener Weg: Remote-Deploy von macOS/Linux nach Unraid - -Das robusteste Setup fuer Unraid ist jetzt das neue Remote-Deploy-Skript. Es -wird auf deinem Mac oder Linux-Rechner gestartet, verbindet sich per SSH mit -Unraid und fuehrt die eigentliche Installation dort aus: - -```bash -bash docker/deploy-to-unraid.sh \ - --unraid-host 192.168.178.30 \ - --paperless-url http://192.168.178.20:8000 \ - --paperless-token PAPERLESS_TOKEN \ - --ai-api-key OPENAI_KEY \ - --ai-model gpt-4.1-mini -``` - -Das Remote-Skript: - -- verbindet sich per SSH zu Unraid -- kopiert den Host-Installer auf den Server -- uebertraegt optional eine lokale `config.yaml` -- fuehrt die eigentliche Installation direkt auf Unraid aus - -### Direkte Ausfuehrung auf dem Unraid-Server - -Wenn du bereits eine Shell direkt auf Unraid offen hast, kannst du stattdessen -das bereits heruntergeladene Host-Skript oder ein Repo-Checkout dort lokal -ausfuehren: +- Das Repository ist im Broker registriert. +- Die Stack-Quelle zeigt auf einen vollständigen Commit-SHA. +- `PAPERLESS_KIPLUS_TOKEN` ist im Broker-Secret-Store vorhanden. +- Das bestehende Appdata-Verzeichnis `/mnt/user/appdata/paperless-kiplus` + bleibt erhalten. -```bash -bash /boot/config/custom/paperless-kiplus/install-unraid-worker.sh \ - --paperless-url http://192.168.178.20:8000 \ - --paperless-token PAPERLESS_TOKEN \ - --ai-api-key OPENAI_KEY \ - --ai-model gpt-4.1-mini -``` - -oder: - -```bash -bash /pfad/zum/repo/docker/install-unraid-worker.sh \ - --paperless-url http://192.168.178.20:8000 \ - --paperless-token PAPERLESS_TOKEN \ - --ai-api-key OPENAI_KEY \ - --ai-model gpt-4.1-mini -``` +Sicherer Ablauf: -Das Host-Skript: +1. `stack_source_status` +2. `stack_validate` +3. `deploy_plan` +4. `approval_request`, falls erforderlich +5. `deploy_apply` +6. `deployment_status`, `docker_list` und `logs_tail` -- erkennt Docker Compose -- legt das Appdata-Verzeichnis an -- erzeugt oder uebernimmt `config.yaml` -- sichert bestehende Dateien vor dem Ueberschreiben -- schreibt einen Compose-Stack fuer GHCR -- startet oder aktualisiert den Container -- prueft `/api/status` per Health-Check - -Typische Zusatzoptionen: - -```bash ---ssh-user root ---ssh-port 22 ---keep-remote-files true ---data-dir /mnt/user/appdata/paperless-kiplus-worker ---worker-token MEIN_API_TOKEN ---enable-tax-enrichment true ---tax-ai-api-key dummy ---tax-ai-model qwen2.5:7b ---tax-ai-base-url http://192.168.178.30:11434/v1 -``` - -### Alternativ: Unraid-Template - -1. `docker/unraid-template.xml` als eigenes Template importieren. -2. Ein persistentes Appdata-Verzeichnis fuer `/data` angeben. -3. Container starten. -4. Entweder in der Weboberflaeche die YAML einfuegen oder `config.yaml` unter `/data/config/config.yaml` ablegen. -5. Anschliessend ueber die Weboberflaeche `Run`, `Resume`, `Restart` oder `Backfill` starten. +Das Compose referenziert das Secret ausschließlich als +`secret://PAPERLESS_KIPLUS_TOKEN`; der Broker injiziert es erst beim Apply. Das +Worker-Image wird lokal aus dem brokergebundenen Git-Checkout gebaut. Der +Dockerfile pinnt Basisimage und Python-Abhängigkeiten; die Compose-Build-Args +halten den erfolgreich geprüften App-Commit und die App-Version fest. Damit +benötigt die Produktion keinen privaten Registry-Pull. ## Welche Datei ist die produktive Konfiguration? @@ -166,14 +91,16 @@ Wichtig: - `GET /api/logs/download` -> kompletter Log als Text - `GET /api/config/export` -> aktuelle Worker-Konfiguration als JSON-Payload - `GET /api/config/download` -> aktuelle Worker-YAML als Download -- `GET /api/review/entities` -> Dopplungskandidaten und gespeicherte KI-Regeln +- `GET /api/jobs/` -> persistenter Jobstatus +- `POST /api/review/entities/jobs` -> asynchroner Dopplungsscan (`202`) +- `GET /api/review/entities` -> `405`, veralteter blockierender Zugriff - `GET /api/review/rules` -> gespeicherte KI-Regeln fuer Entity-Zuordnungen - `POST /api/review/rules` -> Alias-/Ziel-Regel oder "kein Duplikat" speichern -- `POST /api/review/merge` -> Merge planen oder anwenden +- `POST /api/review/merge` -> Merge asynchron planen oder anwenden (`202`) - `POST /api/config/import` -> neue YAML speichern -- `POST /api/run` -> neuen Lauf starten -- `POST /api/resume` -> pausierten Lauf fortsetzen -- `POST /api/restart` -> frischen Neustart machen +- `POST /api/run` -> neuen Lauf als Job starten (`202`) +- `POST /api/resume` -> pausierten Lauf als Job fortsetzen (`202`) +- `POST /api/restart` -> kontrollierten Neustart als Job starten (`202`) - `POST /api/stop` -> sicher pausieren - `POST /api/stop_now` -> sofort stoppen @@ -186,7 +113,9 @@ Korrespondenten nicht erneut an. ## Debugging ### Container laeuft nicht an -- `docker logs paperless-kiplus-worker` +- In Produktion `deployment_status`, `docker_list` und `logs_tail` im Broker + prüfen. Lokal darf `docker compose logs paperless-kiplus-worker` verwendet + werden. - Pruefe, ob `/data/config/config.yaml` gueltiges YAML ist. - Pruefe, ob `paperless_url`, `paperless_token`, `ai_api_key` und `ai_model` gesetzt sind. - Der produktive Broker-Stack startet den Worker bewusst ohne Root-Rechte als @@ -203,3 +132,14 @@ Korrespondenten nicht erneut an. - Existiert `/data/state/run_state.json`? - Wurde der Lauf mit `stop` pausiert oder durch Provider-Wartezeit angehalten? - Bei `stop_now` ist Resume nur ab dem letzten gespeicherten Fortschritt moeglich. + +### Job bleibt nach einem Worker-Restart stehen + +- `GET /api/jobs/` mit dem Worker-Token prüfen. +- Read-only-Scans werden einmal sicher fortgesetzt. +- Schreibjobs erhalten absichtlich `interrupted`; Paperless-Zustand prüfen und + erst danach kontrolliert erneut auslösen. +- Mit `request_id` in den redigierten Broker-Logs suchen. + +Details zu Idempotenz, Parallelität, Aufbewahrung und Rollback stehen unter +[Cloudflare-sichere Langläufer](./cloudflare-long-running-jobs.md). diff --git a/requirements.txt b/requirements.txt index fb8f4c1..8b9cfeb 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,2 +1,6 @@ -requests>=2.34.2 -PyYAML>=6.0.3 +# Purpose: Exact runtime dependency lock for CI, GHCR, and Broker-local builds. +# Input/Output: pip resolves only these reviewed versions into the worker image. +# Invariants: Update versions deliberately and verify tests plus Docker build. +# Debugging: Run `python3 -m pip install -r requirements.txt` in a fresh venv. +requests==2.34.2 +PyYAML==6.0.3 diff --git a/src/background_jobs.py b/src/background_jobs.py new file mode 100644 index 0000000..590478f --- /dev/null +++ b/src/background_jobs.py @@ -0,0 +1,484 @@ +"""Persist Cloudflare-safe background-job metadata in SQLite. + +Purpose: +- Admit potentially long HTTP operations quickly and run them outside the + request thread. +- Keep job/request IDs, exact progress, terminal results, and safe errors + available across browser reloads and worker restarts. + +Input / Output: +- Input: validated operation names, minimal non-secret parameters, optional + idempotency keys, and bounded Python runner callbacks. +- Output: public job dictionaries that never expose callbacks, secret values, + raw exceptions, database paths, or idempotency hashes. + +Important invariants: +- SQLite transactions serialize active resource ownership across HTTP threads. +- Mutating jobs are never retried automatically after process interruption. +- Only explicitly allowlisted read-only jobs may be resubmitted on startup. +- Terminal rows expire after 24 hours and are capped at 200 records. + +How to debug: +- Set ``LOG_LEVEL=DEBUG`` and search logs for ``job_id`` plus ``request_id``. +- Inspect ``state/background_jobs.sqlite3`` locally with SQLite tooling; never + copy raw job databases into issues or chat messages. +""" + +from __future__ import annotations + +import hashlib +import json +import logging +import re +import sqlite3 +import uuid +from collections.abc import Callable, Iterator +from concurrent.futures import ThreadPoolExecutor +from contextlib import contextmanager +from datetime import UTC, datetime, timedelta +from pathlib import Path +from typing import Any + +LOGGER = logging.getLogger("paperless_worker.jobs") +TERMINAL_STATUSES = {"succeeded", "failed", "interrupted", "cancelled"} +ACTIVE_STATUSES = {"queued", "running"} +RETENTION_HOURS = 24 +MAX_TERMINAL_JOBS = 200 +MAX_IDEMPOTENCY_KEY_CHARS = 200 + +ProgressCallback = Callable[[dict[str, Any]], None] +JobRunner = Callable[[ProgressCallback], dict[str, Any]] + + +class JobConflictError(RuntimeError): + """Signal that one serialized resource already has an active owner.""" + + def __init__(self, existing_job: dict[str, Any]) -> None: + super().__init__("Für diese Ressource läuft bereits ein Hintergrundjob.") + self.existing_job = existing_job + + +def _utc_now() -> str: + return datetime.now(UTC).isoformat() + + +def _safe_user_error(exc: Exception) -> tuple[str, str]: + """Map exceptions to bounded messages without echoing provider details.""" + + if isinstance(exc, ValueError): + return ( + "invalid_input", + ( + "Die Job-Ausführung wurde wegen ungültiger Eingaben abgebrochen. " + "Mit request_id in den redigierten Worker-Logs nachsehen." + ), + ) + return ( + "operation_failed", + "Der Hintergrundjob ist fehlgeschlagen. Mit request_id in den redigierten Worker-Logs nachsehen.", + ) + + +class PersistentJobStore: + """Thread-safe, process-persistent job admission and execution store.""" + + def __init__(self, database_path: Path, *, max_workers: int = 3) -> None: + self.database_path = database_path + self.database_path.parent.mkdir(parents=True, exist_ok=True) + self._executor = ThreadPoolExecutor( + max_workers=max(1, min(int(max_workers), 8)), + thread_name_prefix="paperless-job", + ) + self._initialize_schema() + try: + self.database_path.chmod(0o600) + except OSError: + LOGGER.warning("Job-Datenbankrechte konnten nicht auf 0600 gesetzt werden.") + + def _connect(self) -> sqlite3.Connection: + connection = sqlite3.connect( + self.database_path, + timeout=15, + isolation_level=None, + ) + connection.row_factory = sqlite3.Row + connection.execute("PRAGMA busy_timeout = 15000") + connection.execute("PRAGMA journal_mode = WAL") + connection.execute("PRAGMA foreign_keys = ON") + return connection + + @contextmanager + def _connection(self) -> Iterator[sqlite3.Connection]: + """Yield one short-lived connection and always close its file handle.""" + + connection = self._connect() + try: + yield connection + finally: + connection.close() + + def _initialize_schema(self) -> None: + # Why this exists: CREATE TABLE IF NOT EXISTS is an additive migration + # that leaves existing worker state and sorter files untouched. + with self._connection() as connection: + connection.executescript( + """ + CREATE TABLE IF NOT EXISTS background_jobs ( + job_id TEXT PRIMARY KEY, + request_id TEXT NOT NULL UNIQUE, + operation TEXT NOT NULL, + resource_key TEXT, + idempotency_hash TEXT, + status TEXT NOT NULL, + phase TEXT NOT NULL, + params_json TEXT NOT NULL, + progress_json TEXT, + result_json TEXT, + error_code TEXT, + error_message TEXT, + attempt_count INTEGER NOT NULL DEFAULT 0, + max_attempts INTEGER NOT NULL DEFAULT 1, + created_at TEXT NOT NULL, + started_at TEXT, + finished_at TEXT, + updated_at TEXT NOT NULL + ); + CREATE UNIQUE INDEX IF NOT EXISTS uq_background_jobs_idempotency + ON background_jobs(operation, idempotency_hash) + WHERE idempotency_hash IS NOT NULL; + CREATE UNIQUE INDEX IF NOT EXISTS uq_background_jobs_active_resource + ON background_jobs(resource_key) + WHERE resource_key IS NOT NULL + AND status IN ('queued', 'running'); + CREATE INDEX IF NOT EXISTS ix_background_jobs_updated_at + ON background_jobs(updated_at); + """ + ) + + @staticmethod + def _idempotency_hash(operation: str, idempotency_key: str | None) -> str | None: + raw = str(idempotency_key or "").strip() + if not raw: + return None + if len(raw) > MAX_IDEMPOTENCY_KEY_CHARS: + raise ValueError( + f"Idempotency-Key darf höchstens {MAX_IDEMPOTENCY_KEY_CHARS} Zeichen lang sein." + ) + return hashlib.sha256(f"{operation}\0{raw}".encode()).hexdigest() + + def _cleanup(self, connection: sqlite3.Connection) -> None: + cutoff = (datetime.now(UTC) - timedelta(hours=RETENTION_HOURS)).isoformat() + connection.execute( + "DELETE FROM background_jobs WHERE status IN ('succeeded', 'failed', 'interrupted', 'cancelled') AND updated_at < ?", + (cutoff,), + ) + connection.execute( + """ + DELETE FROM background_jobs + WHERE job_id IN ( + SELECT job_id FROM background_jobs + WHERE status IN ('succeeded', 'failed', 'interrupted', 'cancelled') + ORDER BY updated_at DESC + LIMIT -1 OFFSET ? + ) + """, + (MAX_TERMINAL_JOBS,), + ) + + def submit( + self, + *, + operation: str, + params: dict[str, Any], + resource_key: str | None, + idempotency_key: str | None, + runner: JobRunner, + replace_active_operations: set[str] | None = None, + ) -> tuple[dict[str, Any], bool]: + """Atomically admit one job and start it in the bounded executor. + + Example: two simultaneous ``sorter_run`` requests receive one active + owner; a repeated request with the same idempotency key receives the + original job with ``deduplicated=True``. + """ + + normalized_operation = re.sub(r"[^a-z0-9_]+", "_", str(operation).lower()).strip("_") + if not normalized_operation: + raise ValueError("Hintergrundjob benötigt einen gültigen Operationstyp.") + idempotency_hash = self._idempotency_hash(normalized_operation, idempotency_key) + job_id = f"job_{uuid.uuid4().hex}" + request_id = f"req_{uuid.uuid4().hex}" + now = _utc_now() + params_text = json.dumps(params, ensure_ascii=False, separators=(",", ":"), sort_keys=True) + + with self._connection() as connection: + connection.execute("BEGIN IMMEDIATE") + try: + self._cleanup(connection) + if idempotency_hash: + existing = connection.execute( + "SELECT * FROM background_jobs WHERE operation = ? AND idempotency_hash = ?", + (normalized_operation, idempotency_hash), + ).fetchone() + if existing is not None: + connection.commit() + return self._public_row(existing), True + if resource_key: + active = connection.execute( + "SELECT * FROM background_jobs WHERE resource_key = ? AND status IN ('queued', 'running') ORDER BY created_at LIMIT 1", + (resource_key,), + ).fetchone() + if active is not None: + replaceable = replace_active_operations or set() + if str(active["operation"]) not in replaceable: + connection.commit() + raise JobConflictError(self._public_row(active)) + # Why this exists: a restart must be able to supersede + # an active sorter run without opening a second write + # lane. The old callback may finish later, but all of + # its terminal UPDATEs require status='running' and can + # therefore never overwrite this cancellation. + connection.execute( + """ + UPDATE background_jobs + SET status = 'cancelled', phase = 'superseded', + error_code = 'superseded_by_restart', + error_message = 'Der Lauf wurde durch einen kontrollierten Restart ersetzt.', + finished_at = ?, updated_at = ? + WHERE job_id = ? AND status IN ('queued', 'running') + """, + (now, now, active["job_id"]), + ) + connection.execute( + """ + INSERT INTO background_jobs ( + job_id, request_id, operation, resource_key, + idempotency_hash, status, phase, params_json, + attempt_count, max_attempts, created_at, updated_at + ) VALUES (?, ?, ?, ?, ?, 'queued', 'queued', ?, 0, 1, ?, ?) + """, + ( + job_id, + request_id, + normalized_operation, + resource_key, + idempotency_hash, + params_text, + now, + now, + ), + ) + connection.commit() + except Exception: + connection.rollback() + raise + + self._executor.submit(self._execute, job_id, request_id, runner) + job = self.get(job_id) + if job is None: # pragma: no cover - defensive database invariant + raise RuntimeError("Der eben angelegte Hintergrundjob ist nicht lesbar.") + return job, False + + def close(self, *, wait: bool = True) -> None: + """Stop accepting work and optionally wait for active callbacks.""" + + self._executor.shutdown(wait=wait, cancel_futures=False) + + def _execute(self, job_id: str, request_id: str, runner: JobRunner) -> None: + now = _utc_now() + with self._connection() as connection: + changed = connection.execute( + """ + UPDATE background_jobs + SET status = 'running', phase = 'running', attempt_count = attempt_count + 1, + started_at = COALESCE(started_at, ?), updated_at = ? + WHERE job_id = ? AND status = 'queued' + """, + (now, now, job_id), + ).rowcount + if changed != 1: + return + + def progress_callback(progress: dict[str, Any]) -> None: + self.update_progress(job_id, progress) + + try: + result = runner(progress_callback) + safe_result = result if isinstance(result, dict) else {"ok": True} + finished_at = _utc_now() + with self._connection() as connection: + connection.execute( + """ + UPDATE background_jobs + SET status = 'succeeded', phase = 'completed', result_json = ?, + finished_at = ?, updated_at = ? + WHERE job_id = ? AND status = 'running' + """, + ( + json.dumps(safe_result, ensure_ascii=False, separators=(",", ":")), + finished_at, + finished_at, + job_id, + ), + ) + except Exception as exc: # noqa: BLE001 + error_code, error_message = _safe_user_error(exc) + finished_at = _utc_now() + LOGGER.error( + "background_job_failed job_id=%s request_id=%s operation_error=%s", + job_id, + request_id, + type(exc).__name__, + ) + with self._connection() as connection: + connection.execute( + """ + UPDATE background_jobs + SET status = 'failed', phase = 'failed', error_code = ?, + error_message = ?, finished_at = ?, updated_at = ? + WHERE job_id = ? AND status = 'running' + """, + (error_code, error_message, finished_at, finished_at, job_id), + ) + + def update_progress(self, job_id: str, progress: dict[str, Any]) -> None: + """Persist exact counters only; callers use ``None`` when ETA is unknown.""" + + allowed = { + "phase", + "progress_percent", + "estimated_seconds_remaining", + "total", + "completed", + "scanned", + "updated", + "skipped", + "failed", + } + safe_progress = {key: progress.get(key) for key in allowed if key in progress} + now = _utc_now() + with self._connection() as connection: + connection.execute( + "UPDATE background_jobs SET progress_json = ?, phase = ?, updated_at = ? WHERE job_id = ? AND status = 'running'", + ( + json.dumps(safe_progress, ensure_ascii=False, separators=(",", ":")), + str(safe_progress.get("phase") or "running")[:80], + now, + job_id, + ), + ) + + def reconcile_startup( + self, + *, + read_only_runners: dict[str, Callable[[dict[str, Any]], JobRunner]] | None = None, + ) -> dict[str, int]: + """Recover read-only queued work and safely interrupt all mutations.""" + + factories = read_only_runners or {} + resumed = 0 + interrupted = 0 + with self._connection() as connection: + rows = connection.execute( + "SELECT * FROM background_jobs WHERE status IN ('queued', 'running') ORDER BY created_at" + ).fetchall() + for row in rows: + operation = str(row["operation"]) + if operation in factories: + connection.execute( + "UPDATE background_jobs SET status = 'queued', phase = 'recovered', updated_at = ? WHERE job_id = ?", + (_utc_now(), row["job_id"]), + ) + try: + params = json.loads(row["params_json"] or "{}") + except json.JSONDecodeError: + params = {} + self._executor.submit( + self._execute, + row["job_id"], + row["request_id"], + factories[operation](params), + ) + resumed += 1 + continue + now = _utc_now() + connection.execute( + """ + UPDATE background_jobs + SET status = 'interrupted', phase = 'interrupted', + error_code = 'worker_restarted', + error_message = 'Der Worker wurde während des Jobs neu gestartet. Zustand prüfen und die Aktion kontrolliert erneut auslösen.', + finished_at = ?, updated_at = ? + WHERE job_id = ? + """, + (now, now, row["job_id"]), + ) + interrupted += 1 + return {"resumed_read_only": resumed, "interrupted_mutations": interrupted} + + def get(self, job_id: str) -> dict[str, Any] | None: + if not re.fullmatch(r"job_[0-9a-f]{32}", str(job_id or "")): + return None + with self._connection() as connection: + row = connection.execute( + "SELECT * FROM background_jobs WHERE job_id = ?", + (job_id,), + ).fetchone() + cutoff = (datetime.now(UTC) - timedelta(hours=RETENTION_HOURS)).isoformat() + if ( + row is not None + and row["status"] in TERMINAL_STATUSES + and str(row["updated_at"]) < cutoff + ): + connection.execute( + "DELETE FROM background_jobs WHERE job_id = ?", + (job_id,), + ) + row = None + return self._public_row(row) if row is not None else None + + @staticmethod + def _public_row(row: sqlite3.Row) -> dict[str, Any]: + progress = json.loads(row["progress_json"]) if row["progress_json"] else None + result = json.loads(row["result_json"]) if row["result_json"] else None + payload: dict[str, Any] = { + "job_id": row["job_id"], + "request_id": row["request_id"], + "operation": row["operation"], + "status": row["status"], + "phase": row["phase"], + "progress": progress, + "progress_percent": progress.get("progress_percent") if progress else None, + "estimated_seconds_remaining": ( + progress.get("estimated_seconds_remaining") if progress else None + ), + "attempt_count": row["attempt_count"], + "max_attempts": row["max_attempts"], + "created_at": row["created_at"], + "started_at": row["started_at"], + "finished_at": row["finished_at"], + "updated_at": row["updated_at"], + } + if row["status"] == "succeeded": + payload["result"] = result or {} + if row["status"] in {"failed", "interrupted", "cancelled"}: + payload["error"] = { + "code": row["error_code"] or "operation_failed", + "message": row["error_message"] or "Der Hintergrundjob ist fehlgeschlagen.", + "request_id": row["request_id"], + } + return payload + + +def admission_payload(job: dict[str, Any], *, deduplicated: bool) -> dict[str, Any]: + """Build the stable HTTP 202 response shared by all long operations.""" + + job_id = str(job["job_id"]) + return { + "ok": True, + "status": job["status"], + "job_id": job_id, + "request_id": job["request_id"], + "status_url": f"/api/jobs/{job_id}", + "deduplicated": bool(deduplicated), + } diff --git a/src/worker_api.py b/src/worker_api.py index 3901e41..480667f 100644 --- a/src/worker_api.py +++ b/src/worker_api.py @@ -28,24 +28,34 @@ from __future__ import annotations import argparse -from dataclasses import dataclass -from datetime import UTC, datetime, timedelta import json import logging import os -from pathlib import Path +import re +import secrets import shlex import subprocess import sys import threading import time +import uuid +from collections.abc import Callable +from dataclasses import dataclass +from datetime import UTC, datetime, timedelta from http import HTTPStatus from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer -from typing import Any, Optional -from urllib.parse import parse_qs, urlparse +from pathlib import Path +from typing import Any +from urllib.parse import urlparse import yaml +from background_jobs import ( + JobConflictError, + PersistentJobStore, + ProgressCallback, + admission_payload, +) from entity_review import ( EntityRecord, build_ai_prompt_context, @@ -57,8 +67,6 @@ ) from paperless_ai_sorter import ( RUN_PAUSE_EXIT_CODE, - RUN_STATE_FILE_DEFAULT, - STOP_REQUEST_FILE_DEFAULT, ConfigError, PaperlessApiError, PaperlessClient, @@ -71,6 +79,30 @@ FORCE_STOP_GRACE_SECONDS = 5.0 DEFAULT_PORT = 8787 DEFAULT_HOST = "0.0.0.0" +APP_VERSION = ( + str(os.getenv("PAPERLESS_KIPLUS_APP_VERSION", "dev")).strip()[:40] or "dev" +) +_raw_image_commit = str( + os.getenv("PAPERLESS_KIPLUS_IMAGE_COMMIT", "unknown") +).strip().lower() +APP_COMMIT = ( + _raw_image_commit + if re.fullmatch(r"[0-9a-f]{40}", _raw_image_commit) + else "unknown" +) +_SENSITIVE_LOG_PATTERNS = ( + re.compile(r"(?i)(authorization\s*[:=]\s*bearer\s+)[^\s,;]+"), + re.compile(r"(?i)((?:api[_-]?key|token|secret|password|cookie)\s*[:=]\s*)[^\s,;]+"), +) + + +def redact_worker_text(value: Any) -> str: + """Mask common credentials before text reaches files, memory, or the UI.""" + + text = str(value or "") + for pattern in _SENSITIVE_LOG_PATTERNS: + text = pattern.sub(r"\1[REDACTED]", text) + return text WORKER_WEB_UI_HTML = """ @@ -330,10 +362,19 @@ const maxDocumentsInput = document.getElementById('max-documents'); const dryRunSelect = document.getElementById('dry-run'); - const savedToken = localStorage.getItem('paperless_kiplus_worker_token') || ''; + // Why sessionStorage: bearer tokens must survive a page reload, but they + // should not remain on disk after the browser tab is closed. + const legacyToken = localStorage.getItem('paperless_kiplus_worker_token') || ''; + const savedToken = sessionStorage.getItem('paperless_kiplus_worker_token') || legacyToken; + localStorage.removeItem('paperless_kiplus_worker_token'); tokenInput.value = savedToken; tokenInput.addEventListener('change', () => { - localStorage.setItem('paperless_kiplus_worker_token', tokenInput.value.trim()); + const token = tokenInput.value.trim(); + if (token) { + sessionStorage.setItem('paperless_kiplus_worker_token', token); + } else { + sessionStorage.removeItem('paperless_kiplus_worker_token'); + } }); function apiHeaders() { @@ -345,8 +386,8 @@ return headers; } - async function apiJson(path, method = 'GET', payload = null) { - const options = { method, headers: apiHeaders() }; + async function apiJson(path, method = 'GET', payload = null, extraHeaders = {}) { + const options = { method, headers: { ...apiHeaders(), ...extraHeaders } }; if (payload !== null) { options.headers['Content-Type'] = 'application/json'; options.body = JSON.stringify(payload); @@ -358,11 +399,91 @@ data = JSON.parse(text); } if (!response.ok) { - throw new Error(data.message || text || `HTTP ${response.status}`); + const error = new Error(data.message || text || `HTTP ${response.status}`); + error.status = response.status; + throw error; } return data; } + const activeJobsKey = 'paperless_kiplus_active_jobs'; + + function storedJobs() { + try { + const parsed = JSON.parse(localStorage.getItem(activeJobsKey) || '[]'); + return Array.isArray(parsed) ? parsed : []; + } catch (_error) { + return []; + } + } + + function rememberJob(job) { + const jobs = storedJobs().filter(item => item.job_id !== job.job_id); + jobs.push({ job_id: job.job_id, request_id: job.request_id, status_url: job.status_url }); + localStorage.setItem(activeJobsKey, JSON.stringify(jobs.slice(-10))); + } + + function forgetJob(jobId) { + localStorage.setItem( + activeJobsKey, + JSON.stringify(storedJobs().filter(item => item.job_id !== jobId)) + ); + } + + function wait(milliseconds) { + return new Promise(resolve => setTimeout(resolve, milliseconds)); + } + + async function monitorJob(admission) { + rememberJob(admission); + let delayMs = 750; + let transientFailures = 0; + while (true) { + try { + const job = await apiJson(admission.status_url); + transientFailures = 0; + const percent = job.progress_percent; + const progressText = percent === null || percent === undefined + ? `Phase: ${job.phase || job.status}; Fortschritt noch nicht bestimmbar.` + : `Fortschritt: ${Number(percent).toFixed(2)}%.`; + showMessage(`${progressText} Request-ID: ${job.request_id}`); + await refreshStatus(); + if (['succeeded', 'failed', 'interrupted', 'cancelled'].includes(job.status)) { + forgetJob(job.job_id); + if (job.status !== 'succeeded') { + const error = job.error || {}; + throw new Error(`${error.message || 'Job fehlgeschlagen.'} Request-ID: ${error.request_id || job.request_id}`); + } + const result = job.result || {}; + showMessage(`${result.message || 'Hintergrundjob abgeschlossen.'} Request-ID: ${job.request_id}`); + await refreshLog(); + return result; + } + delayMs = Math.min(10000, Math.round(delayMs * 1.7)); + } catch (error) { + if (String(error.message || '').includes('Request-ID:')) { + showMessage(error.message, true); + throw error; + } + if (error.status === 404) { + forgetJob(admission.job_id); + const message = `Jobstatus ist abgelaufen oder nicht mehr vorhanden. Request-ID: ${admission.request_id}`; + showMessage(message, true); + throw new Error(message); + } + transientFailures += 1; + if (transientFailures >= 12) { + const message = `Jobstatus nach 12 Versuchen nicht erreichbar. Der Job bleibt für einen späteren Reload gespeichert. Request-ID: ${admission.request_id}`; + showMessage(message, true); + throw new Error(message); + } + showMessage(`Jobstatus vorübergehend nicht erreichbar; erneuter Versuch. Request-ID: ${admission.request_id}`, true); + delayMs = Math.min(10000, Math.round(delayMs * 2)); + } + await wait(delayMs); + } + } + function showMessage(text, isError = false) { actionMessage.style.display = 'block'; actionMessage.style.background = isError ? '#fdeeee' : '#fff6ea'; @@ -435,7 +556,21 @@ async function callAction(path, payload = {}) { try { - const data = await apiJson(path, 'POST', payload); + const isLongOperation = ['/api/run', '/api/resume', '/api/restart'].includes(path); + const idempotencyKey = isLongOperation && globalThis.crypto?.randomUUID + ? globalThis.crypto.randomUUID() + : `${Date.now()}-${Math.random()}`; + const data = await apiJson( + path, + 'POST', + payload, + isLongOperation ? { 'Idempotency-Key': idempotencyKey } : {} + ); + if (data.job_id && data.status_url) { + showMessage(`Job angenommen. Request-ID: ${data.request_id}`); + monitorJob(data).catch(() => {}); + return; + } if (data.status) { renderStatus(data.status); } @@ -489,6 +624,9 @@ refreshConfig(); tick(); + for (const job of storedJobs()) { + monitorJob(job).catch(() => {}); + } setInterval(tick, 5000); @@ -802,15 +940,15 @@ function apiHeaders() { const headers = { 'Accept': 'application/json' }; - const token = localStorage.getItem('paperless_kiplus_worker_token') || ''; + const token = sessionStorage.getItem('paperless_kiplus_worker_token') || ''; if (token.trim()) { headers.Authorization = `Bearer ${token.trim()}`; } return headers; } - async function apiJson(path, method = 'GET', payload = null) { - const options = { method, headers: apiHeaders() }; + async function apiJson(path, method = 'GET', payload = null, extraHeaders = {}) { + const options = { method, headers: { ...apiHeaders(), ...extraHeaders } }; if (payload !== null) { options.headers['Content-Type'] = 'application/json'; options.body = JSON.stringify(payload); @@ -819,11 +957,105 @@ const text = await response.text(); const data = text.trim() ? JSON.parse(text) : {}; if (!response.ok) { - throw new Error(data.message || text || `HTTP ${response.status}`); + const error = new Error(data.message || text || `HTTP ${response.status}`); + error.status = response.status; + throw error; } return data; } + const reviewJobsKey = 'paperless_kiplus_review_jobs'; + + function storedReviewJobs() { + try { + const parsed = JSON.parse(localStorage.getItem(reviewJobsKey) || '[]'); + return Array.isArray(parsed) ? parsed : []; + } catch (_error) { + return []; + } + } + + function rememberReviewJob(job, purpose) { + const jobs = storedReviewJobs().filter(item => item.job_id !== job.job_id); + jobs.push({ + job_id: job.job_id, + request_id: job.request_id, + status_url: job.status_url, + purpose + }); + localStorage.setItem(reviewJobsKey, JSON.stringify(jobs.slice(-10))); + } + + function forgetReviewJob(jobId) { + localStorage.setItem( + reviewJobsKey, + JSON.stringify(storedReviewJobs().filter(item => item.job_id !== jobId)) + ); + } + + function wait(milliseconds) { + return new Promise(resolve => setTimeout(resolve, milliseconds)); + } + + async function monitorReviewJob(admission, purpose) { + rememberReviewJob(admission, purpose); + let delayMs = 750; + let transientFailures = 0; + while (true) { + try { + const job = await apiJson(admission.status_url); + transientFailures = 0; + const percent = job.progress_percent; + const progressText = percent === null || percent === undefined + ? `Phase: ${job.phase || job.status}; Fortschritt noch nicht bestimmbar.` + : `Fortschritt: ${Number(percent).toFixed(2)}%.`; + showMessage(listMessage, `${progressText} Request-ID: ${job.request_id}`); + if (['succeeded', 'failed', 'interrupted', 'cancelled'].includes(job.status)) { + forgetReviewJob(job.job_id); + if (job.status !== 'succeeded') { + const error = job.error || {}; + throw new Error(`${error.message || 'Job fehlgeschlagen.'} Request-ID: ${error.request_id || job.request_id}`); + } + return job.result || {}; + } + delayMs = Math.min(10000, Math.round(delayMs * 1.7)); + } catch (error) { + if (String(error.message || '').includes('Request-ID:')) { + throw error; + } + if (error.status === 404) { + forgetReviewJob(admission.job_id); + throw new Error(`Jobstatus ist abgelaufen oder nicht mehr vorhanden. Request-ID: ${admission.request_id}`); + } + transientFailures += 1; + if (transientFailures >= 12) { + throw new Error(`Jobstatus nach 12 Versuchen nicht erreichbar. Der Job bleibt für einen späteren Reload gespeichert. Request-ID: ${admission.request_id}`); + } + showMessage( + listMessage, + `Jobstatus vorübergehend nicht erreichbar; erneuter Versuch. Request-ID: ${admission.request_id}`, + true + ); + delayMs = Math.min(10000, Math.round(delayMs * 2)); + } + await wait(delayMs); + } + } + + async function submitReviewJob(path, payload, purpose) { + const idempotencyKey = globalThis.crypto?.randomUUID + ? globalThis.crypto.randomUUID() + : `${Date.now()}-${Math.random()}`; + const admission = await apiJson( + path, + 'POST', + payload, + { 'Idempotency-Key': idempotencyKey } + ); + showMessage(listMessage, `Job angenommen. Request-ID: ${admission.request_id}`); + return monitorReviewJob(admission, purpose); + } + function showMessage(element, text, isError = false) { element.style.display = 'block'; element.classList.toggle('error', Boolean(isError)); @@ -869,7 +1101,10 @@ document.getElementById('stat-rules').textContent = payload.rules.length; document.getElementById('stat-doc-types').textContent = payload.entities.document_type.length; document.getElementById('stat-correspondents').textContent = payload.entities.correspondent.length; - document.getElementById('ai-context').textContent = payload.ai_context || 'Noch keine gespeicherten Prefer-/Merge-Regeln vorhanden.'; + const contextLines = (payload.rules || []) + .filter(rule => ['prefer', 'merge'].includes(rule.action)) + .map(rule => `${rule.alias_name} → ${rule.canonical_name}${rule.context ? `: ${rule.context}` : ''}`); + document.getElementById('ai-context').textContent = contextLines.join('\n') || 'Noch keine gespeicherten Prefer-/Merge-Regeln vorhanden.'; } function candidateMatches(candidate, query, typeFilter) { @@ -961,7 +1196,11 @@ async function refresh() { hideMessage(listMessage); const threshold = Number(document.getElementById('threshold-input').value || 0.84); - const payload = await apiJson(`/api/review/entities?threshold=${encodeURIComponent(threshold)}`); + const payload = await submitReviewJob( + '/api/review/entities/jobs', + { threshold }, + 'scan' + ); state.payload = payload; state.candidates = payload.candidates || []; renderStats(payload); @@ -996,8 +1235,12 @@ document.getElementById('dry-run-btn').addEventListener('click', async () => { try { - const result = await apiJson('/api/review/merge', 'POST', { ...buildDecisionPayload('merge'), dry_run: true }); - showMessage(actionMessage, `Merge-Plan: ${result.affected_count || 0} Dokumente würden umgehängt. Vorschau-IDs: ${(result.affected_document_ids_preview || []).join(', ') || '-'}\n${result.message || ''}`); + const result = await submitReviewJob( + '/api/review/merge', + { ...buildDecisionPayload('merge'), dry_run: true }, + 'merge_plan' + ); + showMessage(actionMessage, `Merge-Plan: ${result.affected_count || 0} Dokumente würden umgehängt.\n${result.message || ''}`); } catch (error) { showMessage(actionMessage, error.message, true); } @@ -1010,7 +1253,11 @@ if (!ok) { return; } - const result = await apiJson('/api/review/merge', 'POST', { ...payload, dry_run: false }); + const result = await submitReviewJob( + '/api/review/merge', + { ...payload, dry_run: false }, + 'merge_apply' + ); showMessage(actionMessage, result.message || `Merge angewendet: ${result.updated_count || 0} Dokumente aktualisiert.`); await refresh(); } catch (error) { @@ -1028,7 +1275,23 @@ } }); - refresh().catch(error => showMessage(listMessage, `Review konnte nicht geladen werden: ${error.message}`, true)); + const pendingReviewJob = storedReviewJobs().at(-1); + if (pendingReviewJob) { + monitorReviewJob(pendingReviewJob, pendingReviewJob.purpose) + .then(payload => { + if (pendingReviewJob.purpose === 'scan') { + state.payload = payload; + state.candidates = payload.candidates || []; + renderStats(payload); + renderCandidates(); + } else { + showMessage(actionMessage, payload.message || 'Hintergrundjob abgeschlossen.'); + } + }) + .catch(error => showMessage(listMessage, error.message, true)); + } else { + refresh().catch(error => showMessage(listMessage, `Review konnte nicht geladen werden: ${error.message}`, true)); + } @@ -1059,6 +1322,7 @@ class WorkerPaths: metrics_file: Path entity_review_rules_file: Path worker_meta_file: Path + jobs_file: Path log_file: Path @@ -1084,6 +1348,7 @@ def __init__( metrics_file=data_dir / "state" / "run_metrics.json", entity_review_rules_file=data_dir / "state" / "entity_review_rules.json", worker_meta_file=data_dir / "state" / "worker_meta.json", + jobs_file=data_dir / "state" / "background_jobs.sqlite3", log_file=data_dir / "logs" / "worker.log", ) self.sorter_command = list(sorter_command) @@ -1156,6 +1421,15 @@ def __init__( self._ensure_directories() self._restore_state_on_startup() + self.jobs = PersistentJobStore(self.paths.jobs_file, max_workers=3) + recovery = self.jobs.reconcile_startup( + read_only_runners={"review_scan": self._recovered_review_scan_runner}, + ) + LOGGER.info( + "background_jobs_reconciled resumed_read_only=%s interrupted_mutations=%s", + recovery["resumed_read_only"], + recovery["interrupted_mutations"], + ) def _ensure_directories(self) -> None: for path in ( @@ -1201,7 +1475,7 @@ def _safe_float(value: Any, default: float = 0.0) -> float: return default def _append_log_line(self, stream_name: str, line: str) -> None: - stripped = line.rstrip("\n") + stripped = redact_worker_text(line.rstrip("\n")) prefixed = f"[{stream_name}] {stripped}" self.log_lines.append(prefixed) self.log_lines = self.log_lines[-1500:] @@ -1329,7 +1603,7 @@ def _refresh_config_state(self) -> None: self.config_validation_message = "Konfiguration ist gültig." except ConfigError as exc: self.config_validation_ok = False - self.config_validation_message = str(exc) + self.config_validation_message = redact_worker_text(str(exc))[:500] def import_config_yaml(self, yaml_text: str, *, source: str) -> dict[str, Any]: raw = str(yaml_text or "") @@ -1338,7 +1612,7 @@ def import_config_yaml(self, yaml_text: str, *, source: str) -> dict[str, Any]: except yaml.YAMLError as exc: raise ValueError(f"YAML ist ungültig: {exc}") from exc if not isinstance(parsed, dict): - raise ValueError("Die Worker-Konfiguration muss ein YAML-Objekt sein.") + raise TypeError("Die Worker-Konfiguration muss ein YAML-Objekt sein.") self.paths.config_file.parent.mkdir(parents=True, exist_ok=True) self.paths.config_file.write_text(raw, encoding="utf-8") self.config_source = source @@ -1518,8 +1792,8 @@ def _stream_reader(self, pipe: Any, *, stream_name: str) -> None: finally: try: pipe.close() - except Exception: - pass + except (OSError, ValueError) as exc: + LOGGER.debug("Prozess-Stream konnte nicht geschlossen werden: %s", type(exc).__name__) def _persist_force_stop_resume_state(self) -> None: if self.paths.run_state_file.exists(): @@ -1563,7 +1837,7 @@ def _resume() -> None: try: self.resume_run(force=True) except Exception as exc: # noqa: BLE001 - LOGGER.error("Auto-Resume fehlgeschlagen: %s", exc) + LOGGER.error("Auto-Resume fehlgeschlagen: %s", type(exc).__name__) self.auto_resume_timer = threading.Timer(delay, _resume) self.auto_resume_timer.daemon = True @@ -1705,10 +1979,14 @@ def start_run( resume_run: bool = False, ) -> dict[str, Any]: with self.lock: - if self.running and not force: - self.last_status = "skipped_running" - self.last_message = "run skipped because another run is active" - return self.status_payload() + if self.running: + if not force: + self.last_status = "skipped_running" + self.last_message = "run skipped because another run is active" + return self.status_payload() + raise ValueError( + "Ein Lauf ist bereits aktiv. Für einen kontrollierten Wechsel /api/restart verwenden." + ) self._start_process( dry_run=dry_run, all_documents=all_documents, @@ -1769,6 +2047,9 @@ def restart_run( force: bool = True, backfill_existing_documents: bool | None = None, ) -> dict[str, Any]: + # Why the lock is released while waiting: process finalization happens + # in another thread and must acquire the same lock to mark ``running`` + # false. Holding it here would turn every restart into a 45s deadlock. with self.lock: base_payload: dict[str, Any] = {} if self.paths.run_state_file.exists(): @@ -1782,19 +2063,23 @@ def restart_run( restart_backfill = bool(mode.get("backfill_existing_documents", False)) if backfill_existing_documents is not None: restart_backfill = bool(backfill_existing_documents) + was_running = self.running + + if was_running: + self.force_stop() + deadline = time.time() + 45.0 + while self.running and time.time() < deadline: + time.sleep(0.25) if self.running: - self.force_stop() - deadline = time.time() + 45.0 - while self.running and time.time() < deadline: - time.sleep(0.25) - if self.running: - raise RuntimeError("Vorheriger Prozess konnte nicht rechtzeitig beendet werden.") + raise RuntimeError("Vorheriger Prozess konnte nicht rechtzeitig beendet werden.") + + with self.lock: self._clear_restart_state_files() self.latest_runtime_payload = {} self.resume_available = False self.pause_reason = "" self.auto_resume_at = None - return self.start_run(force=force, backfill_existing_documents=restart_backfill) + return self.start_run(force=force, backfill_existing_documents=restart_backfill) def reset_metrics(self) -> dict[str, Any]: with self.lock: @@ -1821,7 +2106,11 @@ def reset_failed_documents(self) -> dict[str, Any]: path.unlink() deleted_count += 1 except OSError as exc: - LOGGER.warning("Konnte Failed-Datei nicht löschen (%s): %s", path, exc) + LOGGER.warning( + "Konnte Failed-Datei nicht löschen (%s): %s", + path.name, + type(exc).__name__, + ) self._refresh_failed_state_counts() self.last_status = "failed_docs_reset" self.last_message = f"failed/quarantine documents reset ({deleted_count} files)" @@ -2035,7 +2324,12 @@ def _document_ids_for_entity( continue return sorted(set(document_ids)) - def merge_review_entities(self, payload: dict[str, Any]) -> dict[str, Any]: + def merge_review_entities( + self, + payload: dict[str, Any], + *, + progress_callback: ProgressCallback | None = None, + ) -> dict[str, Any]: entity_type = str(payload.get("entity_type") or "").strip() alias_id = int(payload.get("alias_id")) canonical_id = int(payload.get("canonical_id")) @@ -2047,12 +2341,31 @@ def merge_review_entities(self, payload: dict[str, Any]) -> dict[str, Any]: delete_alias = bool(payload.get("delete_alias", True)) field = self._document_field_for_entity(entity_type) endpoint = self._entity_endpoint(entity_type) + if progress_callback: + progress_callback( + { + "phase": "discovering_documents", + "progress_percent": None, + "estimated_seconds_remaining": None, + "completed": 0, + } + ) client = self._load_paperless_client() document_ids = self._document_ids_for_entity( client, entity_type=entity_type, entity_id=alias_id, ) + if progress_callback: + progress_callback( + { + "phase": "planning" if dry_run else "updating_documents", + "progress_percent": 100.0 if dry_run else (0.0 if document_ids else 100.0), + "estimated_seconds_remaining": None, + "total": len(document_ids), + "completed": 0, + } + ) if dry_run: return { @@ -2060,7 +2373,6 @@ def merge_review_entities(self, payload: dict[str, Any]) -> dict[str, Any]: "dry_run": True, "message": "Dry-Run abgeschlossen. Es wurden keine Paperless-Daten geändert.", "affected_count": len(document_ids), - "affected_document_ids_preview": document_ids[:25], "delete_alias": delete_alias, } @@ -2068,6 +2380,17 @@ def merge_review_entities(self, payload: dict[str, Any]) -> dict[str, Any]: for document_id in document_ids: client.update_document(document_id, {field: canonical_id}) updated_count += 1 + if progress_callback: + progress_callback( + { + "phase": "updating_documents", + "progress_percent": round(updated_count * 100 / len(document_ids), 2), + "estimated_seconds_remaining": None, + "total": len(document_ids), + "completed": updated_count, + "updated": updated_count, + } + ) deleted_alias = False delete_warning = "" @@ -2076,9 +2399,15 @@ def merge_review_entities(self, payload: dict[str, Any]) -> dict[str, Any]: client._request("DELETE", f"{endpoint}{alias_id}/", retries=2) deleted_alias = True except PaperlessApiError as exc: + LOGGER.warning( + "review_alias_delete_failed entity_type=%s alias_id=%s error_type=%s", + entity_type, + alias_id, + type(exc).__name__, + ) delete_warning = ( "Dokumente wurden umgehängt, aber der alte Paperless-Eintrag " - f"konnte nicht gelöscht werden: {exc}" + "konnte nicht gelöscht werden. request_id im Jobstatus für die Logs verwenden." ) rule = upsert_review_rule( @@ -2098,13 +2427,241 @@ def merge_review_entities(self, payload: dict[str, Any]) -> dict[str, Any]: "affected_count": len(document_ids), "updated_count": updated_count, "deleted_alias": deleted_alias, - "affected_document_ids_preview": document_ids[:25], "rule": rule.to_payload(), "ai_context": build_ai_prompt_context(self._load_review_rules()), } + def _sorter_progress_payload(self) -> dict[str, Any]: + """Return exact sorter counters without document titles or URLs.""" + + with self.lock: + percent = ( + round(self.progress_percent, 2) + if self.progress_total_documents > 0 + else None + ) + return { + "phase": "running" if self.running else str(self.last_status or "completed"), + "progress_percent": percent, + "estimated_seconds_remaining": None, + "total": self.progress_total_documents, + "completed": self.progress_completed_documents, + "scanned": self.progress_scanned, + "updated": self.progress_updated, + "skipped": self.progress_skipped, + "failed": self.progress_failed, + } + + def _sorter_result_payload(self) -> dict[str, Any]: + """Return the minimal persistent result needed after a page reload.""" + + with self.lock: + return { + "status": self.last_status, + "message": self.last_message, + "last_exit_code": self.last_exit_code, + "resume_available": self.resume_available, + "scanned": self.progress_scanned, + "updated": self.progress_updated, + "skipped": self.progress_skipped, + "failed": self.progress_failed, + } + + def _run_sorter_background_job( + self, + operation: str, + params: dict[str, Any], + progress_callback: ProgressCallback, + ) -> dict[str, Any]: + """Start one sorter mode, then follow the real subprocess to terminal state.""" + + if operation == "sorter_run": + self.start_run( + force=bool(params.get("force", False)), + dry_run=bool(params.get("dry_run", False)), + all_documents=bool(params.get("all_documents", False)), + max_documents=int(params.get("max_documents", 0)), + backfill_existing_documents=bool( + params.get("backfill_existing_documents", False) + ), + ) + elif operation == "sorter_resume": + self.resume_run(force=bool(params.get("force", False))) + elif operation == "sorter_restart": + self.restart_run( + force=bool(params.get("force", True)), + backfill_existing_documents=params.get("backfill_existing_documents"), + ) + else: # pragma: no cover - caller allowlist invariant + raise ValueError(f"Nicht unterstützte Sorter-Operation: {operation}") + + while True: + progress_callback(self._sorter_progress_payload()) + with self.lock: + running = self.running + if not running: + break + time.sleep(1.0) + result = self._sorter_result_payload() + if result.get("status") == "error": + raise RuntimeError("Der Sorter-Lauf ist fehlgeschlagen.") + return result + + @staticmethod + def _safe_review_scan_result(payload: dict[str, Any]) -> dict[str, Any]: + """Persist only review data required to redraw the authenticated page.""" + + return { + "ok": True, + "threshold": payload.get("threshold"), + "entities": payload.get("entities") or {}, + "candidates": payload.get("candidates") or [], + "rules": payload.get("rules") or [], + } + + @staticmethod + def _safe_review_merge_result(payload: dict[str, Any]) -> dict[str, Any]: + """Omit document IDs, paths, raw provider errors, and generated AI context.""" + + return { + key: payload.get(key) + for key in ( + "ok", + "dry_run", + "message", + "affected_count", + "updated_count", + "deleted_alias", + "delete_alias", + ) + if key in payload + } + + def _review_scan_runner( + self, + threshold: float, + ) -> Callable[[ProgressCallback], dict[str, Any]]: + def runner(progress_callback: ProgressCallback) -> dict[str, Any]: + progress_callback( + { + "phase": "loading_entities", + "progress_percent": None, + "estimated_seconds_remaining": None, + } + ) + result = self.entity_review_payload(threshold=threshold) + return self._safe_review_scan_result(result) + + return runner + + def _recovered_review_scan_runner( + self, + params: dict[str, Any], + ) -> Callable[[ProgressCallback], dict[str, Any]]: + """Rebuild only the allowlisted read-only scan after a worker restart.""" + + threshold = self._safe_float(params.get("threshold"), 0.84) + return self._review_scan_runner(threshold) + + def submit_background_job( + self, + operation: str, + payload: dict[str, Any], + *, + idempotency_key: str | None, + ) -> tuple[dict[str, Any], bool]: + """Validate one public operation and atomically submit its safe runner.""" + + if operation in {"sorter_run", "sorter_resume", "sorter_restart"}: + max_documents = int(payload.get("max_documents", 0) or 0) + if max_documents < 0 or max_documents > 1_000_000: + raise ValueError("max_documents muss zwischen 0 und 1000000 liegen.") + params = { + "force": bool(payload.get("force", operation == "sorter_restart")), + "dry_run": bool(payload.get("dry_run", False)), + "all_documents": bool(payload.get("all_documents", False)), + "max_documents": max_documents, + "backfill_existing_documents": payload.get( + "backfill_existing_documents", False + ), + } + return self.jobs.submit( + operation=operation, + params=params, + resource_key="paperless_write", + idempotency_key=idempotency_key, + runner=lambda progress: self._run_sorter_background_job( + operation, + params, + progress, + ), + replace_active_operations=( + {"sorter_run", "sorter_resume"} + if operation == "sorter_restart" + else None + ), + ) + + if operation == "review_scan": + threshold = self._safe_float(payload.get("threshold"), 0.84) + if threshold < 0.5 or threshold > 1.0: + raise ValueError("threshold muss zwischen 0.5 und 1.0 liegen.") + return self.jobs.submit( + operation=operation, + params={"threshold": threshold}, + resource_key="review_scan", + idempotency_key=idempotency_key, + runner=self._review_scan_runner(threshold), + ) + + if operation == "review_merge": + entity_type = str(payload.get("entity_type") or "").strip() + self._entity_endpoint(entity_type) + alias_id = int(payload.get("alias_id")) + canonical_id = int(payload.get("canonical_id")) + if alias_id <= 0 or canonical_id <= 0 or alias_id == canonical_id: + raise ValueError("Alias und Ziel benötigen unterschiedliche positive IDs.") + safe_payload = { + "entity_type": entity_type, + "alias_id": alias_id, + "alias_name": str(payload.get("alias_name") or "")[:200], + "canonical_id": canonical_id, + "canonical_name": str(payload.get("canonical_name") or "")[:200], + "context": str(payload.get("context") or "")[:2000], + "dry_run": bool(payload.get("dry_run", True)), + "delete_alias": bool(payload.get("delete_alias", True)), + } + + def merge_runner(progress: ProgressCallback) -> dict[str, Any]: + result = self.merge_review_entities( + safe_payload, + progress_callback=progress, + ) + return self._safe_review_merge_result(result) + + # Names/context are intentionally kept only in this in-memory + # closure. Interrupted writes are never replayed after restart. + persisted_params = { + "entity_type": entity_type, + "alias_id": alias_id, + "canonical_id": canonical_id, + "dry_run": safe_payload["dry_run"], + "delete_alias": safe_payload["delete_alias"], + } + return self.jobs.submit( + operation=operation, + params=persisted_params, + resource_key="paperless_write", + idempotency_key=idempotency_key, + runner=merge_runner, + ) + + raise ValueError(f"Nicht unterstützte Hintergrundoperation: {operation}") + def status_payload(self) -> dict[str, Any]: return { + "app_version": APP_VERSION, + "app_commit": APP_COMMIT, "status": self.last_status, "message": self.last_message, "running": self.running, @@ -2197,6 +2754,8 @@ def heimdall_payload(self) -> dict[str, Any]: last_run_at = self.last_finished or self.last_started last_run_value = last_run_at.isoformat() if last_run_at else "never" details = [ + f"Version: {APP_VERSION}", + f"Image commit: {APP_COMMIT}", f"Config: {config_state}", f"Last run: {last_run_value}", ] @@ -2210,6 +2769,7 @@ def heimdall_payload(self) -> dict[str, Any]: {"label": "Status", "value": worker_state}, {"label": "Running", "value": "yes" if self.running else "no"}, {"label": "Failed", "value": str(self.last_failed)}, + {"label": "Version", "value": APP_VERSION}, ], "details": details, } @@ -2226,6 +2786,9 @@ def _json_response(self, payload: dict[str, Any], status: HTTPStatus = HTTPStatu encoded = json.dumps(payload, ensure_ascii=False).encode("utf-8") self.send_response(status) self.send_header("Content-Type", "application/json; charset=utf-8") + self.send_header("Cache-Control", "no-store") + self.send_header("X-Content-Type-Options", "nosniff") + self.send_header("Referrer-Policy", "no-referrer") self.send_header("Content-Length", str(len(encoded))) self.end_headers() self.wfile.write(encoded) @@ -2240,6 +2803,11 @@ def _text_response( encoded = text.encode("utf-8") self.send_response(status) self.send_header("Content-Type", content_type) + # Config and log downloads may contain private operational data. They + # must not become browser or Cloudflare cache entries. + self.send_header("Cache-Control", "no-store") + self.send_header("X-Content-Type-Options", "nosniff") + self.send_header("Referrer-Policy", "no-referrer") self.send_header("Content-Length", str(len(encoded))) self.end_headers() self.wfile.write(encoded) @@ -2248,19 +2816,34 @@ def _read_json_body(self) -> dict[str, Any]: content_length = int(self.headers.get("Content-Length", "0") or 0) if content_length <= 0: return {} + if content_length > 1024 * 1024: + raise ValueError("JSON-Body darf höchstens 1 MiB groß sein.") raw = self.rfile.read(content_length).decode("utf-8") if not raw.strip(): return {} payload = json.loads(raw) if not isinstance(payload, dict): - raise ValueError("JSON-Body muss ein Objekt sein.") + raise TypeError("JSON-Body muss ein Objekt sein.") return payload def _is_authorized(self) -> bool: if not self.manager.auth_token: return True header = str(self.headers.get("Authorization") or "").strip() - return header == f"Bearer {self.manager.auth_token}" + return secrets.compare_digest(header, f"Bearer {self.manager.auth_token}") + + def _admit_job(self, operation: str, payload: dict[str, Any]) -> None: + """Return the stable 202 admission contract for one long operation.""" + + job, deduplicated = self.manager.submit_background_job( + operation, + payload, + idempotency_key=self.headers.get("Idempotency-Key"), + ) + self._json_response( + admission_payload(job, deduplicated=deduplicated), + status=HTTPStatus.ACCEPTED, + ) def _require_auth(self) -> bool: if self._is_authorized(): @@ -2271,9 +2854,10 @@ def _require_auth(self) -> bool: ) return False - def do_GET(self) -> None: # noqa: N802 + def do_GET(self) -> None: parsed = urlparse(self.path) path = parsed.path + request_id = f"req_{uuid.uuid4().hex}" try: if path == "/": self._text_response(WORKER_WEB_UI_HTML, content_type="text/html; charset=utf-8") @@ -2292,6 +2876,16 @@ def do_GET(self) -> None: # noqa: N802 if path == "/api/status": self._json_response(self.manager.status_payload()) return + if path.startswith("/api/jobs/"): + job = self.manager.jobs.get(path.removeprefix("/api/jobs/")) + if job is None: + self._json_response( + {"ok": False, "message": "Job nicht gefunden.", "request_id": request_id}, + status=HTTPStatus.NOT_FOUND, + ) + return + self._json_response({"ok": True, **job}) + return if path == "/api/logs": self._json_response( { @@ -2320,25 +2914,39 @@ def do_GET(self) -> None: # noqa: N802 ) return if path == "/api/review/entities": - query = parse_qs(parsed.query) - threshold_raw = (query.get("threshold") or ["0.84"])[0] - try: - threshold = float(threshold_raw) - except (TypeError, ValueError): - threshold = 0.84 - self._json_response(self.manager.entity_review_payload(threshold=threshold)) + self._json_response( + { + "ok": False, + "message": "Live-Review ist asynchron. POST /api/review/entities/jobs verwenden.", + "request_id": request_id, + }, + status=HTTPStatus.METHOD_NOT_ALLOWED, + ) return if path == "/api/review/rules": self._json_response(self.manager.entity_review_rules_payload()) return self._json_response({"ok": False, "message": "Not found"}, status=HTTPStatus.NOT_FOUND) except Exception as exc: # noqa: BLE001 - LOGGER.exception("GET %s fehlgeschlagen: %s", path, exc) - self._json_response({"ok": False, "message": str(exc)}, status=HTTPStatus.INTERNAL_SERVER_ERROR) + LOGGER.error( + "worker_http_get_failed path=%s request_id=%s error_type=%s", + path, + request_id, + type(exc).__name__, + ) + self._json_response( + { + "ok": False, + "message": "Interner Fehler. Mit request_id in den redigierten Logs nachsehen.", + "request_id": request_id, + }, + status=HTTPStatus.INTERNAL_SERVER_ERROR, + ) - def do_POST(self) -> None: # noqa: N802 + def do_POST(self) -> None: parsed = urlparse(self.path) path = parsed.path + request_id = f"req_{uuid.uuid4().hex}" try: if not path.startswith("/api/"): self._json_response({"ok": False, "message": "Not found"}, status=HTTPStatus.NOT_FOUND) @@ -2347,25 +2955,16 @@ def do_POST(self) -> None: # noqa: N802 return payload = self._read_json_body() if path == "/api/run": - status_payload = self.manager.start_run( - force=bool(payload.get("force", False)), - dry_run=bool(payload.get("dry_run", False)), - all_documents=bool(payload.get("all_documents", False)), - max_documents=int(payload.get("max_documents", 0) or 0), - backfill_existing_documents=bool(payload.get("backfill_existing_documents", False)), - ) - self._json_response({"ok": True, "message": "Lauf gestartet.", "status": status_payload}) + self._admit_job("sorter_run", payload) return if path == "/api/resume": - status_payload = self.manager.resume_run(force=bool(payload.get("force", True))) - self._json_response({"ok": True, "message": "Resume ausgelöst.", "status": status_payload}) + self._admit_job("sorter_resume", payload) return if path == "/api/restart": - status_payload = self.manager.restart_run( - force=bool(payload.get("force", True)), - backfill_existing_documents=payload.get("backfill_existing_documents"), - ) - self._json_response({"ok": True, "message": "Neustart ausgelöst.", "status": status_payload}) + self._admit_job("sorter_restart", payload) + return + if path == "/api/review/entities/jobs": + self._admit_job("review_scan", payload) return if path == "/api/stop": status_payload = self.manager.request_stop() @@ -2395,17 +2994,50 @@ def do_POST(self) -> None: # noqa: N802 self._json_response(result) return if path == "/api/review/merge": - result = self.manager.merge_review_entities(payload) - self._json_response(result) + self._admit_job("review_merge", payload) return self._json_response({"ok": False, "message": "Not found"}, status=HTTPStatus.NOT_FOUND) - except ValueError as exc: - self._json_response({"ok": False, "message": str(exc)}, status=HTTPStatus.BAD_REQUEST) + except (TypeError, ValueError) as exc: + self._json_response( + { + "ok": False, + "message": redact_worker_text(str(exc))[:500], + "request_id": request_id, + }, + status=HTTPStatus.BAD_REQUEST, + ) + except JobConflictError as exc: + existing = exc.existing_job + self._json_response( + { + "ok": False, + "message": str(exc), + "request_id": request_id, + "active_job": { + "job_id": existing["job_id"], + "request_id": existing["request_id"], + "status_url": f"/api/jobs/{existing['job_id']}", + }, + }, + status=HTTPStatus.CONFLICT, + ) except Exception as exc: # noqa: BLE001 - LOGGER.exception("POST %s fehlgeschlagen: %s", path, exc) - self._json_response({"ok": False, "message": str(exc)}, status=HTTPStatus.INTERNAL_SERVER_ERROR) + LOGGER.error( + "worker_http_post_failed path=%s request_id=%s error_type=%s", + path, + request_id, + type(exc).__name__, + ) + self._json_response( + { + "ok": False, + "message": "Interner Fehler. Mit request_id in den redigierten Logs nachsehen.", + "request_id": request_id, + }, + status=HTTPStatus.INTERNAL_SERVER_ERROR, + ) - def log_message(self, format: str, *args: Any) -> None: # noqa: A003 + def log_message(self, format: str, *args: Any) -> None: LOGGER.info("worker_http | %s", format % args) @@ -2439,8 +3071,14 @@ def parse_args() -> argparse.Namespace: def main() -> int: args = parse_args() + log_level_name = str(os.getenv("PAPERLESS_KIPLUS_LOG_LEVEL", "INFO")).strip().upper() + log_level = getattr(logging, log_level_name, None) + if not isinstance(log_level, int): + raise TypeError( + "PAPERLESS_KIPLUS_LOG_LEVEL muss DEBUG, INFO, WARNING, ERROR oder CRITICAL sein." + ) logging.basicConfig( - level=logging.INFO, + level=log_level, format="%(asctime)s | %(levelname)s | %(name)s | %(message)s", ) manager = WorkerManager( @@ -2450,7 +3088,9 @@ def main() -> int: ) server = WorkerHttpServer((args.host, int(args.port)), manager) LOGGER.info( - "Starte Paperless KIplus Worker | host=%s port=%s data_dir=%s", + "Starte Paperless KIplus Worker | version=%s image_commit=%s host=%s port=%s data_dir=%s", + APP_VERSION, + APP_COMMIT, args.host, args.port, manager.paths.data_dir, @@ -2461,6 +3101,7 @@ def main() -> int: LOGGER.info("Worker wird beendet ...") finally: server.server_close() + manager.jobs.close(wait=False) return 0 diff --git a/tests/test_background_jobs.py b/tests/test_background_jobs.py new file mode 100644 index 0000000..5971c49 --- /dev/null +++ b/tests/test_background_jobs.py @@ -0,0 +1,315 @@ +"""Tests for the persistent Cloudflare-safe background-job store. + +Purpose: +- Prove that admission is idempotent and parallel-safe. +- Protect restart recovery and safe terminal error behavior. + +Input / Output: +- Input: temporary SQLite files and deterministic in-process job callbacks. +- Output: assertions against the same public job dictionaries served by HTTP. + +Important invariants: +- At most one active job owns a serialized resource. +- Mutations become interrupted after restart; allowlisted reads may resume. +- Secret-like exception text never reaches persisted public errors. + +How to debug: +- Run `python3 -m unittest tests.test_background_jobs -v`. +- Add a temporary `print(store.get(job_id))` only in a local checkout; the + temporary databases are deleted at test completion. +""" + +from __future__ import annotations + +import sys +import tempfile +import threading +import time +import unittest +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +SRC_DIR = ROOT / "src" +if str(SRC_DIR) not in sys.path: + sys.path.insert(0, str(SRC_DIR)) + +from background_jobs import JobConflictError, PersistentJobStore + + +def _wait_for_status( + store: PersistentJobStore, + job_id: str, + statuses: set[str], + *, + timeout: float = 5.0, +) -> dict[str, object]: + """Poll a local store with a short deterministic test deadline.""" + + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + job = store.get(job_id) + if job and job["status"] in statuses: + return job + time.sleep(0.01) + raise AssertionError(f"Job {job_id} erreichte {sorted(statuses)} nicht rechtzeitig.") + + +class PersistentJobStoreTests(unittest.TestCase): + """Covers storage, parallel admission, recovery, and redaction.""" + + def test_success_progress_and_idempotent_replay(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + store = PersistentJobStore(Path(tmp_dir) / "jobs.sqlite3") + try: + def runner(progress): + progress( + { + "phase": "counting", + "progress_percent": 50.0, + "estimated_seconds_remaining": None, + "completed": 1, + "total": 2, + } + ) + return {"ok": True, "message": "fertig"} + + first, first_deduplicated = store.submit( + operation="review_scan", + params={"threshold": 0.84}, + resource_key="review_scan", + idempotency_key="same-browser-action", + runner=runner, + ) + terminal = _wait_for_status(store, first["job_id"], {"succeeded"}) + second, second_deduplicated = store.submit( + operation="review_scan", + params={"threshold": 0.84}, + resource_key="review_scan", + idempotency_key="same-browser-action", + runner=runner, + ) + + self.assertFalse(first_deduplicated) + self.assertTrue(second_deduplicated) + self.assertEqual(first["job_id"], second["job_id"]) + self.assertEqual(terminal["result"]["message"], "fertig") + self.assertIsNone(terminal["estimated_seconds_remaining"]) + finally: + store.close() + + def test_parallel_admission_keeps_one_resource_owner(self) -> None: + # Five fresh databases catch timing-sensitive regressions instead of + # proving the unique-index path only once. + for repetition in range(5): + with self.subTest(repetition=repetition), tempfile.TemporaryDirectory() as tmp_dir: + store = PersistentJobStore(Path(tmp_dir) / "jobs.sqlite3") + release = threading.Event() + try: + def blocking_runner(_progress, event=release): + event.wait(5) + return {"ok": True} + + first, _ = store.submit( + operation="review_merge", + params={"alias_id": 1, "canonical_id": 2}, + resource_key="paperless_write", + idempotency_key="first", + runner=blocking_runner, + ) + _wait_for_status(store, first["job_id"], {"running"}) + + conflicts: list[str] = [] + unexpected: list[Exception] = [] + + def contend( + index: int, + current_store=store, + current_conflicts=conflicts, + current_unexpected=unexpected, + ) -> None: + try: + current_store.submit( + operation="sorter_run", + params={}, + resource_key="paperless_write", + idempotency_key=f"contender-{index}", + runner=lambda _progress: {"ok": True}, + ) + except JobConflictError as exc: + current_conflicts.append(exc.existing_job["job_id"]) + except Exception as exc: # noqa: BLE001 - test captures thread failures + current_unexpected.append(exc) + + threads = [ + threading.Thread(target=contend, args=(index,)) + for index in range(8) + ] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=5) + + self.assertEqual(unexpected, []) + self.assertEqual(conflicts, [first["job_id"]] * 8) + finally: + release.set() + store.close() + + def test_restart_supersedes_sorter_but_not_merge(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + store = PersistentJobStore(Path(tmp_dir) / "jobs.sqlite3", max_workers=2) + release = threading.Event() + try: + active, _ = store.submit( + operation="sorter_run", + params={}, + resource_key="paperless_write", + idempotency_key="active-sorter", + runner=lambda _progress: (release.wait(5), {"ok": True})[1], + ) + _wait_for_status(store, active["job_id"], {"running"}) + replacement, _ = store.submit( + operation="sorter_restart", + params={}, + resource_key="paperless_write", + idempotency_key="restart", + runner=lambda _progress: {"ok": True}, + replace_active_operations={"sorter_run", "sorter_resume", "sorter_restart"}, + ) + + cancelled = _wait_for_status(store, active["job_id"], {"cancelled"}) + _wait_for_status(store, replacement["job_id"], {"succeeded"}) + self.assertEqual(cancelled["error"]["code"], "superseded_by_restart") + finally: + release.set() + store.close() + + def test_restart_marks_mutation_interrupted(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + database = Path(tmp_dir) / "jobs.sqlite3" + first_store = PersistentJobStore(database) + release = threading.Event() + second_store: PersistentJobStore | None = None + try: + job, _ = first_store.submit( + operation="review_merge", + params={"alias_id": 1, "canonical_id": 2}, + resource_key="paperless_write", + idempotency_key=None, + runner=lambda _progress: (release.wait(5), {"ok": True})[1], + ) + _wait_for_status(first_store, job["job_id"], {"running"}) + second_store = PersistentJobStore(database) + recovery = second_store.reconcile_startup() + interrupted = _wait_for_status( + second_store, + job["job_id"], + {"interrupted"}, + ) + + self.assertEqual(recovery["interrupted_mutations"], 1) + self.assertEqual(interrupted["error"]["code"], "worker_restarted") + self.assertEqual(interrupted["attempt_count"], 1) + finally: + release.set() + first_store.close() + if second_store is not None: + second_store.close() + + def test_allowlisted_read_is_resumed_once_after_restart(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + database = Path(tmp_dir) / "jobs.sqlite3" + first_store = PersistentJobStore(database) + release = threading.Event() + second_store: PersistentJobStore | None = None + try: + job, _ = first_store.submit( + operation="review_scan", + params={"threshold": 0.9}, + resource_key="review_scan", + idempotency_key=None, + runner=lambda _progress: (release.wait(5), {"old": True})[1], + ) + _wait_for_status(first_store, job["job_id"], {"running"}) + second_store = PersistentJobStore(database) + recovery = second_store.reconcile_startup( + read_only_runners={ + "review_scan": lambda params: ( + lambda _progress: {"threshold": params["threshold"], "resumed": True} + ) + } + ) + terminal = _wait_for_status(second_store, job["job_id"], {"succeeded"}) + + self.assertEqual(recovery["resumed_read_only"], 1) + self.assertTrue(terminal["result"]["resumed"]) + self.assertEqual(terminal["attempt_count"], 2) + finally: + release.set() + first_store.close() + if second_store is not None: + second_store.close() + + def test_failure_does_not_persist_secret_text(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + store = PersistentJobStore(Path(tmp_dir) / "jobs.sqlite3") + try: + def failing_runner(_progress): + raise ValueError("token=do-not-leak provider payload") + + job, _ = store.submit( + operation="review_scan", + params={}, + resource_key="review_scan", + idempotency_key=None, + runner=failing_runner, + ) + terminal = _wait_for_status(store, job["job_id"], {"failed"}) + rendered = str(terminal).lower() + + self.assertNotIn("do-not-leak", rendered) + self.assertEqual(terminal["error"]["code"], "invalid_input") + self.assertEqual(terminal["max_attempts"], 1) + finally: + store.close() + + def test_overlong_idempotency_key_is_rejected(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + store = PersistentJobStore(Path(tmp_dir) / "jobs.sqlite3") + try: + with self.assertRaisesRegex(ValueError, "höchstens 200"): + store.submit( + operation="review_scan", + params={}, + resource_key="review_scan", + idempotency_key="x" * 201, + runner=lambda _progress: {"ok": True}, + ) + finally: + store.close() + + def test_expired_terminal_job_is_no_longer_public(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + store = PersistentJobStore(Path(tmp_dir) / "jobs.sqlite3") + try: + job, _ = store.submit( + operation="review_scan", + params={}, + resource_key="review_scan", + idempotency_key=None, + runner=lambda _progress: {"ok": True}, + ) + _wait_for_status(store, job["job_id"], {"succeeded"}) + with store._connection() as connection: + connection.execute( + "UPDATE background_jobs SET updated_at = '2000-01-01T00:00:00+00:00' WHERE job_id = ?", + (job["job_id"],), + ) + + self.assertIsNone(store.get(job["job_id"])) + finally: + store.close() + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_container_hardening.py b/tests/test_container_hardening.py index e04fda0..d08b2e1 100644 --- a/tests/test_container_hardening.py +++ b/tests/test_container_hardening.py @@ -9,7 +9,6 @@ import yaml - REPOSITORY_ROOT = Path(__file__).resolve().parents[1] DOCKERFILE_PATH = REPOSITORY_ROOT / "docker" / "Dockerfile" BROKER_COMPOSE_PATH = REPOSITORY_ROOT / "docker" / "docker-compose.unraid-broker.yml" @@ -26,7 +25,9 @@ def test_worker_dockerfile_uses_specific_base_and_non_root_user() -> None: from_instructions = [line for line in dockerfile_lines if line.upper().startswith("FROM ")] user_instructions = [line for line in dockerfile_lines if line.upper().startswith("USER ")] - assert from_instructions == ["FROM python:3.12.13-slim-trixie"] + assert from_instructions == [ + "FROM python:3.12.13-slim-trixie@sha256:229a2c5bfa27522db7815ea81f9bed70af17ccb9de9fc7ad142b1877b5830d36" + ] assert user_instructions assert user_instructions[-1].split(maxsplit=1)[1].lower() not in {"root", "0", "0:0"} @@ -48,3 +49,18 @@ def test_broker_compose_uses_stable_review_port() -> None: worker = compose["services"]["paperless-kiplus-worker"] assert worker["ports"] == ["8788:8788"] + + +def test_broker_compose_build_is_commit_bound() -> None: + """Production must not depend on mutable tags or private GHCR pull access.""" + + compose = yaml.safe_load(BROKER_COMPOSE_PATH.read_text(encoding="utf-8")) + worker = compose["services"]["paperless-kiplus-worker"] + + assert worker["image"] == "paperless-kiplus-worker:1.4.21-e155328" + assert worker["build"]["context"] == "." + assert worker["build"]["dockerfile"] == "docker/Dockerfile" + assert worker["build"]["args"] == { + "APP_COMMIT": "e15532882f58a02e2a033d0f79e6b965eaad39a7", + "APP_VERSION": "1.4.21", + } diff --git a/tests/test_remote_runner_jobs.py b/tests/test_remote_runner_jobs.py new file mode 100644 index 0000000..59b6a70 --- /dev/null +++ b/tests/test_remote_runner_jobs.py @@ -0,0 +1,197 @@ +"""Tests for Home Assistant compatibility with HTTP-202 worker jobs. + +Purpose: +- Protect the remote runner's legacy/202 response adapter and terminal polling. +- Verify that HA actions send an idempotency key without a Home Assistant install. + +Input / Output: +- Input: deterministic fake worker responses and minimal HA/aiohttp modules. +- Output: assertions against the runner fields consumed by HA entities. + +Important invariants: +- A queued job keeps HA polling even before the sorter process is visible. +- Terminal job errors must not be overwritten by a later worker-status refresh. +- Existing legacy immediate responses remain supported. + +How to debug: +- Run `python3 -m unittest tests.test_remote_runner_jobs -v`. +- Inspect `_active_job_status_url`, `last_status`, and `last_message` on failure. +""" + +from __future__ import annotations + +import asyncio +import importlib.util +import sys +import types +import unittest +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +PACKAGE_DIR = ROOT / "custom_components" / "paperless_kiplus" + + +def _load_remote_runner_module(): + """Load the runner with only the HA interfaces used during import.""" + + aiohttp_module = types.ModuleType("aiohttp") + aiohttp_module.ClientTimeout = lambda **kwargs: kwargs + sys.modules.setdefault("aiohttp", aiohttp_module) + + homeassistant_module = types.ModuleType("homeassistant") + core_module = types.ModuleType("homeassistant.core") + helpers_module = types.ModuleType("homeassistant.helpers") + aiohttp_client_module = types.ModuleType("homeassistant.helpers.aiohttp_client") + dispatcher_module = types.ModuleType("homeassistant.helpers.dispatcher") + + class _FakeHomeAssistant: + """Minimal type stub required by the module annotations.""" + + core_module.HomeAssistant = _FakeHomeAssistant + aiohttp_client_module.async_get_clientsession = lambda _hass: None + dispatcher_module.async_dispatcher_send = lambda *_args, **_kwargs: None + homeassistant_module.core = core_module + homeassistant_module.helpers = helpers_module + helpers_module.aiohttp_client = aiohttp_client_module + helpers_module.dispatcher = dispatcher_module + + sys.modules.setdefault("homeassistant", homeassistant_module) + sys.modules.setdefault("homeassistant.core", core_module) + sys.modules.setdefault("homeassistant.helpers", helpers_module) + sys.modules.setdefault("homeassistant.helpers.aiohttp_client", aiohttp_client_module) + sys.modules.setdefault("homeassistant.helpers.dispatcher", dispatcher_module) + + custom_components_module = types.ModuleType("custom_components") + package_module = types.ModuleType("custom_components.paperless_kiplus") + package_module.__path__ = [str(PACKAGE_DIR)] + const_module = types.ModuleType("custom_components.paperless_kiplus.const") + const_module.SIGNAL_STATUS_UPDATED = "paperless_kiplus_test_signal" + sys.modules.setdefault("custom_components", custom_components_module) + sys.modules.setdefault("custom_components.paperless_kiplus", package_module) + sys.modules.setdefault("custom_components.paperless_kiplus.const", const_module) + + module_name = "custom_components.paperless_kiplus.remote_runner" + spec = importlib.util.spec_from_file_location(module_name, PACKAGE_DIR / "remote_runner.py") + if spec is None or spec.loader is None: + raise RuntimeError("Remote-Runner-Modul konnte nicht geladen werden.") + module = importlib.util.module_from_spec(spec) + sys.modules[module_name] = module + spec.loader.exec_module(module) + return module + + +REMOTE_MODULE = _load_remote_runner_module() +RemotePaperlessRunner = REMOTE_MODULE.RemotePaperlessRunner + + +def _bare_runner(): + """Build only the state needed by the methods under test.""" + + runner = RemotePaperlessRunner.__new__(RemotePaperlessRunner) + runner._active_job_status_url = "" + runner._active_job_request_id = "" + runner._poll_task = None + runner._lock = asyncio.Lock() + runner.running = False + runner.resume_available = False + runner.last_status = "idle" + runner.last_message = "not started" + runner.last_stderr_tail = "" + runner.last_exit_code = None + runner.remote_worker_sync_config = False + runner.managed_config_enabled = False + runner.default_dry_run = False + runner.default_all_documents = False + runner.default_max_documents = 0 + runner._notify = lambda: None + return runner + + +class RemoteRunnerJobTests(unittest.TestCase): + """Covers 202 admission, polling, failures, and legacy compatibility.""" + + def test_202_admission_keeps_runner_active(self) -> None: + runner = _bare_runner() + + runner._apply_action_response( + { + "status": "queued", + "job_id": "job_123", + "request_id": "req_123", + "status_url": "/api/jobs/job_123", + } + ) + + self.assertTrue(runner.running) + self.assertEqual(runner.last_status, "queued") + self.assertEqual(runner._active_job_status_url, "/api/jobs/job_123") + self.assertIn("req_123", runner.last_message) + + def test_legacy_immediate_response_remains_supported(self) -> None: + runner = _bare_runner() + applied: list[dict[str, object]] = [] + runner._apply_status_payload = applied.append + + runner._apply_action_response( + {"ok": True, "status": {"running": True, "status": "running"}} + ) + + self.assertEqual(applied, [{"running": True, "status": "running"}]) + + def test_terminal_failure_survives_worker_status_refresh(self) -> None: + runner = _bare_runner() + runner.running = True + runner._active_job_status_url = "/api/jobs/job_failed" + runner._active_job_request_id = "req_failed" + + async def api_json(_method, _path): + return { + "status": "failed", + "error": {"message": "Sichere Jobmeldung.", "request_id": "req_failed"}, + } + + async def refresh_status(): + runner.running = False + runner.last_status = "idle" + runner.last_message = "worker idle" + + runner._api_json = api_json + runner._refresh_status = refresh_status + + asyncio.run(runner._poll_loop()) + + self.assertEqual(runner.last_status, "remote_job_failed") + self.assertEqual(runner.last_message, "Sichere Jobmeldung.") + self.assertEqual(runner._active_job_status_url, "") + + def test_async_run_sends_idempotency_key(self) -> None: + runner = _bare_runner() + captured: dict[str, object] = {} + + async def api_json(method, path, *, payload=None, idempotency_key=None): + captured.update( + method=method, + path=path, + payload=payload, + idempotency_key=idempotency_key, + ) + return { + "status": "queued", + "job_id": "job_123", + "request_id": "req_123", + "status_url": "/api/jobs/job_123", + } + + runner._api_json = api_json + runner._ensure_polling = lambda: None + + asyncio.run(runner.async_run(dry_run=True, max_documents=3)) + + self.assertEqual(captured["method"], "POST") + self.assertEqual(captured["path"], "/api/run") + self.assertEqual(captured["payload"]["max_documents"], 3) + self.assertRegex(str(captured["idempotency_key"]), r"^[0-9a-f]{32}$") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_worker_api.py b/tests/test_worker_api.py index 3de4246..7f8d979 100644 --- a/tests/test_worker_api.py +++ b/tests/test_worker_api.py @@ -18,11 +18,13 @@ import tempfile import textwrap import threading +import time import types import unittest -from urllib.request import urlopen +from contextlib import contextmanager from pathlib import Path - +from urllib.error import HTTPError +from urllib.request import Request, urlopen ROOT = Path(__file__).resolve().parents[1] SRC_DIR = ROOT / "src" @@ -75,7 +77,7 @@ def _simple_safe_load(text: str): return payload -def _simple_safe_dump(payload, allow_unicode=True, sort_keys=False): # noqa: ARG001 +def _simple_safe_dump(payload, allow_unicode=True, sort_keys=False): lines = [] items = payload.items() if isinstance(payload, dict) else [] if sort_keys: @@ -158,12 +160,89 @@ def _request_heimdall_payload(manager: WorkerManager) -> dict[str, object]: body = response.read().decode("utf-8") payload = json.loads(body) if not isinstance(payload, dict): - raise AssertionError("Heimdall response must be a JSON object.") + raise TypeError("Heimdall response must be a JSON object.") return payload finally: server.shutdown() server.server_close() thread.join(timeout=5) + manager.jobs.close() + + +@contextmanager +def _running_worker_server(manager: WorkerManager): + """Serve one manager on an ephemeral local port and close all resources.""" + + server = WorkerHttpServer(("127.0.0.1", 0), manager) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + host, port = server.server_address + try: + yield f"http://{host}:{port}" + finally: + server.shutdown() + server.server_close() + thread.join(timeout=5) + manager.jobs.close() + + +def _json_request( + base_url: str, + path: str, + *, + method: str = "GET", + payload: dict[str, object] | None = None, + token: str | None = None, + idempotency_key: str | None = None, +) -> tuple[int, dict[str, object]]: + """Call the local test server and return JSON for success and HTTP errors.""" + + headers = {"Accept": "application/json"} + if token: + headers["Authorization"] = f"Bearer {token}" + if idempotency_key: + headers["Idempotency-Key"] = idempotency_key + body = None + if payload is not None: + headers["Content-Type"] = "application/json" + body = json.dumps(payload).encode("utf-8") + request = Request( + f"{base_url}{path}", + data=body, + headers=headers, + method=method, + ) + try: + with urlopen(request, timeout=5) as response: + return response.status, json.loads(response.read().decode("utf-8")) + except HTTPError as exc: + try: + return exc.code, json.loads(exc.read().decode("utf-8")) + finally: + exc.close() + + +def _wait_for_http_job( + base_url: str, + status_url: str, + *, + token: str, + timeout: float = 8.0, +) -> dict[str, object]: + """Poll one local HTTP job until it reaches a terminal state.""" + + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + status, job = _json_request(base_url, status_url, token=token) + if status == 200 and job.get("status") in { + "succeeded", + "failed", + "interrupted", + "cancelled", + }: + return job + time.sleep(0.02) + raise AssertionError(f"Job unter {status_url} wurde nicht rechtzeitig terminal.") class ConfigExportTests(unittest.TestCase): @@ -243,6 +322,8 @@ def test_import_config_yaml_updates_worker_state(self) -> None: status = manager.status_payload() self.assertTrue(result["config_validation_ok"]) + self.assertIn("app_version", status) + self.assertIn("app_commit", status) self.assertEqual(status["config_source"], "unit_test") self.assertEqual(status["paperless_base_url"], "https://paperless.example") self.assertEqual(status["config_validation_message"], "Konfiguration ist gültig.") @@ -392,5 +473,271 @@ def test_heimdall_endpoint_omits_sensitive_fields_and_raw_logs(self) -> None: self.assertNotIn("stderr_tail", keys) +class WorkerBackgroundJobApiTests(unittest.TestCase): + """Covers the public HTTP-202 contract and browser recovery hooks.""" + + TOKEN = "test-worker-token" + + @staticmethod + def _valid_yaml() -> str: + return textwrap.dedent( + """ + paperless_url: https://paperless.example + paperless_token: local-test-token + ai_api_key: local-test-ai-key + ai_model: gpt-4.1-mini + ai_base_url: https://api.openai.com/v1 + """ + ).strip() + + def _manager(self, data_dir: str, *, sleep_seconds: float = 0.05) -> WorkerManager: + manager = WorkerManager( + data_dir=Path(data_dir), + sorter_command=[ + sys.executable, + "-c", + f"import time; time.sleep({sleep_seconds}); print('done')", + ], + auth_token=self.TOKEN, + ) + manager.import_config_yaml(self._valid_yaml(), source="unit_test") + return manager + + def test_run_returns_202_and_idempotent_status_resource(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + manager = self._manager(tmp_dir) + with _running_worker_server(manager) as base_url: + first_status, first = _json_request( + base_url, + "/api/run", + method="POST", + payload={"dry_run": True, "max_documents": 1}, + token=self.TOKEN, + idempotency_key="browser-action-1", + ) + second_status, second = _json_request( + base_url, + "/api/run", + method="POST", + payload={"dry_run": True, "max_documents": 1}, + token=self.TOKEN, + idempotency_key="browser-action-1", + ) + + self.assertEqual(first_status, 202) + self.assertEqual(second_status, 202) + self.assertEqual(first["job_id"], second["job_id"]) + self.assertTrue(second["deduplicated"]) + self.assertRegex(str(first["job_id"]), r"^job_[0-9a-f]{32}$") + self.assertRegex(str(first["request_id"]), r"^req_[0-9a-f]{32}$") + self.assertEqual(first["status_url"], f"/api/jobs/{first['job_id']}") + + terminal = _wait_for_http_job( + base_url, + str(first["status_url"]), + token=self.TOKEN, + ) + self.assertEqual(terminal["status"], "succeeded") + self.assertIsNone(terminal["estimated_seconds_remaining"]) + self.assertNotIn("paperless_token", json.dumps(terminal).lower()) + + def test_parallel_write_returns_safe_409_conflict(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + manager = self._manager(tmp_dir, sleep_seconds=0.4) + with _running_worker_server(manager) as base_url: + first_status, first = _json_request( + base_url, + "/api/run", + method="POST", + payload={"dry_run": True}, + token=self.TOKEN, + idempotency_key="first-write", + ) + conflict_status, conflict = _json_request( + base_url, + "/api/review/merge", + method="POST", + payload={ + "entity_type": "correspondent", + "alias_id": 1, + "canonical_id": 2, + "dry_run": True, + }, + token=self.TOKEN, + idempotency_key="second-write", + ) + + self.assertEqual(first_status, 202) + self.assertEqual(conflict_status, 409) + self.assertEqual(conflict["active_job"]["job_id"], first["job_id"]) + self.assertNotIn("params", conflict["active_job"]) + self.assertNotIn("token", json.dumps(conflict).lower()) + _wait_for_http_job(base_url, str(first["status_url"]), token=self.TOKEN) + + def test_restart_replaces_active_sorter_without_parallel_processes(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + manager = self._manager(tmp_dir, sleep_seconds=0.35) + with _running_worker_server(manager) as base_url: + _, first = _json_request( + base_url, + "/api/run", + method="POST", + payload={"dry_run": True}, + token=self.TOKEN, + idempotency_key="run-before-restart", + ) + deadline = time.monotonic() + 3 + while not manager.running and time.monotonic() < deadline: + time.sleep(0.01) + self.assertTrue(manager.running) + + restart_status, restart = _json_request( + base_url, + "/api/restart", + method="POST", + payload={"force": True}, + token=self.TOKEN, + idempotency_key="controlled-restart", + ) + cancelled = _wait_for_http_job( + base_url, + str(first["status_url"]), + token=self.TOKEN, + ) + restarted = _wait_for_http_job( + base_url, + str(restart["status_url"]), + token=self.TOKEN, + ) + + self.assertEqual(restart_status, 202) + self.assertEqual(cancelled["status"], "cancelled") + self.assertEqual(restarted["status"], "succeeded") + self.assertLessEqual(manager.jobs._executor._max_workers, 3) + + def test_read_only_review_scan_is_asynchronous_and_safe(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + manager = self._manager(tmp_dir) + manager.entity_review_payload = lambda threshold: { + "threshold": threshold, + "entities": {"correspondent": [], "document_type": []}, + "candidates": [], + "rules": [], + "ai_context": "must-not-be-persisted", + "private_path": "/data/config/config.yaml", + } + with _running_worker_server(manager) as base_url: + status, admission = _json_request( + base_url, + "/api/review/entities/jobs", + method="POST", + payload={"threshold": 0.9}, + token=self.TOKEN, + idempotency_key="review-scan", + ) + terminal = _wait_for_http_job( + base_url, + str(admission["status_url"]), + token=self.TOKEN, + ) + rendered = json.dumps(terminal) + + self.assertEqual(status, 202) + self.assertEqual(terminal["status"], "succeeded") + self.assertEqual(terminal["result"]["threshold"], 0.9) + self.assertNotIn("must-not-be-persisted", rendered) + self.assertNotIn("private_path", rendered) + + def test_job_status_requires_auth_and_legacy_live_get_is_disabled(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + manager = self._manager(tmp_dir) + manager.entity_review_payload = lambda threshold: { + "threshold": threshold, + "entities": {}, + "candidates": [], + "rules": [], + } + with _running_worker_server(manager) as base_url: + _, admission = _json_request( + base_url, + "/api/review/entities/jobs", + method="POST", + payload={"threshold": 0.84}, + token=self.TOKEN, + idempotency_key="auth-check", + ) + unauthorized_status, _ = _json_request( + base_url, + str(admission["status_url"]), + ) + legacy_status, legacy = _json_request( + base_url, + "/api/review/entities", + token=self.TOKEN, + ) + + self.assertEqual(unauthorized_status, 401) + self.assertEqual(legacy_status, 405) + self.assertIn("POST /api/review/entities/jobs", legacy["message"]) + _wait_for_http_job( + base_url, + str(admission["status_url"]), + token=self.TOKEN, + ) + + def test_invalid_idempotency_key_returns_bounded_400(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + manager = self._manager(tmp_dir) + with _running_worker_server(manager) as base_url: + status, payload = _json_request( + base_url, + "/api/run", + method="POST", + payload={}, + token=self.TOKEN, + idempotency_key="x" * 201, + ) + + self.assertEqual(status, 400) + self.assertIn("höchstens 200", payload["message"]) + self.assertRegex(str(payload["request_id"]), r"^req_[0-9a-f]{32}$") + + def test_worker_log_redaction_happens_before_file_and_memory(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + manager = self._manager(tmp_dir) + try: + manager._append_log_line( + "STDERR", + "Authorization: Bearer top-secret token=other-secret cookie=session-secret", + ) + persisted = manager.paths.log_file.read_text(encoding="utf-8") + rendered = "\n".join(manager.log_lines) + manager.last_stderr_tail + persisted + + self.assertNotIn("top-secret", rendered) + self.assertNotIn("other-secret", rendered) + self.assertNotIn("session-secret", rendered) + self.assertGreaterEqual(rendered.count("[REDACTED]"), 3) + finally: + manager.jobs.close() + + def test_embedded_uis_resume_jobs_with_bounded_backoff(self) -> None: + main_ui = worker_api_module.WORKER_WEB_UI_HTML + review_ui = worker_api_module.ENTITY_REVIEW_HTML + + self.assertIn("paperless_kiplus_active_jobs", main_ui) + self.assertIn("sessionStorage.setItem('paperless_kiplus_worker_token'", main_ui) + self.assertIn("localStorage.removeItem('paperless_kiplus_worker_token')", main_ui) + self.assertIn("for (const job of storedJobs())", main_ui) + self.assertIn("Math.min(10000", main_ui) + self.assertIn("transientFailures >= 12", main_ui) + self.assertIn("Fortschritt noch nicht bestimmbar", main_ui) + self.assertIn("paperless_kiplus_review_jobs", review_ui) + self.assertIn("sessionStorage.getItem('paperless_kiplus_worker_token')", review_ui) + self.assertIn("storedReviewJobs().at(-1)", review_ui) + self.assertIn("Math.min(10000", review_ui) + self.assertIn("transientFailures >= 12", review_ui) + self.assertIn("'/api/review/entities/jobs'", review_ui) + + if __name__ == "__main__": unittest.main()