864 lines
31 KiB
Python
864 lines
31 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
paperless_import.py - Stufe 2 der ecoDMS -> Paperless-ngx Migration.
|
|
|
|
Liest ein Manifest aus ecodms_extract.py (Stufe 1) und spielt es ueber die
|
|
REST-API in Paperless ein. Fuehrt ein eigenes SQLite-Journal, ist damit
|
|
wiederaufsetzbar und tranchenuebergreifend auswertbar.
|
|
|
|
# 1. Trockenlauf: nichts wird geschrieben
|
|
./paperless_import.py tranche_03.json --dry-run
|
|
|
|
# 2. Stammdaten anlegen (Tags, Dokumentarten, Custom Fields)
|
|
./paperless_import.py tranche_03.json --setup
|
|
|
|
# 3. Import
|
|
./paperless_import.py tranche_03.json
|
|
|
|
# 4. Zeitstempel-Skript erzeugen (API kann sie nicht setzen)
|
|
./paperless_import.py --emit-timestamps > /tmp/fix_added.py
|
|
sudo docker compose exec -T webserver python manage.py shell < /tmp/fix_added.py
|
|
|
|
# 5. Nach ALLEN Tranchen: Verknuepfungen aufloesen
|
|
./paperless_import.py --links verknuepfungen.csv
|
|
|
|
# 6. Pruefbericht
|
|
./paperless_import.py --verify
|
|
|
|
Erkenntnisse aus der Testinstanz, die hier kodiert sind:
|
|
|
|
* update_version nimmt nur "document" und "version_label" entgegen -
|
|
Zeitstempel muessen nachtraeglich per Django-Shell gesetzt werden.
|
|
* Rechte werden ueber root_document vererbt. Einmal auf die Wurzel
|
|
setzen genuegt, die Versionen erben (403-Test bestaetigt).
|
|
* Die versions-Liste am Dokument ist ABSTEIGEND sortiert.
|
|
* Innerhalb einer Versionskette muss streng seriell gearbeitet werden,
|
|
parallelisiert wird ueber Dokumente hinweg.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import csv
|
|
import json
|
|
import re
|
|
import sqlite3
|
|
import sys
|
|
import time
|
|
from datetime import datetime
|
|
from pathlib import Path
|
|
|
|
try:
|
|
import requests
|
|
except ImportError:
|
|
sys.exit("Benoetigt python3-requests: sudo zypper in python3-requests")
|
|
|
|
try:
|
|
import yaml
|
|
except ImportError:
|
|
sys.exit("Benoetigt python3-PyYAML: sudo zypper in python3-PyYAML")
|
|
|
|
|
|
TASK_TIMEOUT = 600 # Sekunden je Konsumvorgang
|
|
TASK_POLL_INTERVAL = 2
|
|
|
|
|
|
# ------------------------------------------------------------------- Journal
|
|
|
|
SCHEMA = """
|
|
CREATE TABLE IF NOT EXISTS mapping (
|
|
ecodms_docid INTEGER PRIMARY KEY,
|
|
tranche TEXT,
|
|
paperless_id INTEGER,
|
|
versions_total INTEGER,
|
|
versions_done INTEGER DEFAULT 0,
|
|
status TEXT DEFAULT 'pending',
|
|
last_error TEXT,
|
|
updated_at TEXT
|
|
);
|
|
CREATE TABLE IF NOT EXISTS version_times (
|
|
paperless_id INTEGER PRIMARY KEY,
|
|
ecodms_docid INTEGER,
|
|
version INTEGER,
|
|
added TEXT
|
|
);
|
|
CREATE TABLE IF NOT EXISTS links (
|
|
src_docid INTEGER,
|
|
dst_docid INTEGER,
|
|
resolved INTEGER DEFAULT 0,
|
|
PRIMARY KEY (src_docid, dst_docid)
|
|
);
|
|
"""
|
|
|
|
|
|
class Journal:
|
|
def __init__(self, path: Path):
|
|
self.con = sqlite3.connect(path)
|
|
self.con.row_factory = sqlite3.Row
|
|
self.con.executescript(SCHEMA)
|
|
self.con.commit()
|
|
|
|
def state(self, docid):
|
|
r = self.con.execute(
|
|
"SELECT * FROM mapping WHERE ecodms_docid = ?", (docid,)
|
|
).fetchone()
|
|
return dict(r) if r else None
|
|
|
|
def start(self, docid, tranche, total):
|
|
self.con.execute(
|
|
"INSERT OR IGNORE INTO mapping (ecodms_docid, tranche, versions_total) "
|
|
"VALUES (?,?,?)",
|
|
(docid, tranche, total),
|
|
)
|
|
self.con.commit()
|
|
|
|
def update(self, docid, **fields):
|
|
fields["updated_at"] = datetime.now().isoformat(timespec="seconds")
|
|
sets = ", ".join(f"{k} = ?" for k in fields)
|
|
self.con.execute(
|
|
f"UPDATE mapping SET {sets} WHERE ecodms_docid = ?",
|
|
(*fields.values(), docid),
|
|
)
|
|
self.con.commit()
|
|
|
|
def remember_time(self, paperless_id, docid, version, added):
|
|
self.con.execute(
|
|
"INSERT OR REPLACE INTO version_times VALUES (?,?,?,?)",
|
|
(paperless_id, docid, version, added),
|
|
)
|
|
self.con.commit()
|
|
|
|
def add_link(self, src, dst):
|
|
self.con.execute(
|
|
"INSERT OR IGNORE INTO links (src_docid, dst_docid) VALUES (?,?)",
|
|
(src, dst),
|
|
)
|
|
self.con.commit()
|
|
|
|
|
|
# ----------------------------------------------------------------------- API
|
|
|
|
class Paperless:
|
|
def __init__(self, base, token, dry_run=False):
|
|
self.base = base.rstrip("/")
|
|
self.dry_run = dry_run
|
|
self.s = requests.Session()
|
|
self.s.headers["Authorization"] = f"Token {token}"
|
|
|
|
def _url(self, path):
|
|
return f"{self.base}/api/{path.lstrip('/')}"
|
|
|
|
def get(self, path, **params):
|
|
r = self.s.get(self._url(path), params=params, timeout=60)
|
|
self._raise(r)
|
|
return r.json()
|
|
|
|
def get_all(self, path, **params):
|
|
"""Folgt der Paginierung."""
|
|
out, params = [], {**params, "page_size": 200}
|
|
url = self._url(path)
|
|
while url:
|
|
r = self.s.get(url, params=params, timeout=60)
|
|
self._raise(r)
|
|
data = r.json()
|
|
out.extend(data.get("results", []))
|
|
url, params = data.get("next"), None
|
|
return out
|
|
|
|
@staticmethod
|
|
def _raise(r):
|
|
"""
|
|
raise_for_status() zeigt nur den Statuscode. Die Begruendung steht
|
|
im Antwortkoerper - ohne sie ist ein 400 nicht auswertbar.
|
|
"""
|
|
if r.status_code < 400:
|
|
return
|
|
try:
|
|
body = json.dumps(r.json(), ensure_ascii=False)
|
|
except ValueError:
|
|
body = r.text or ""
|
|
raise requests.HTTPError(
|
|
f"HTTP {r.status_code} bei {r.request.method} {r.url}: "
|
|
f"{' '.join(body.split())[:400]}",
|
|
response=r,
|
|
)
|
|
|
|
def post(self, path, **kwargs):
|
|
if self.dry_run:
|
|
return {"dry_run": True}
|
|
r = self.s.post(self._url(path), timeout=300, **kwargs)
|
|
self._raise(r)
|
|
return r.json()
|
|
|
|
def patch(self, path, payload):
|
|
if self.dry_run:
|
|
return {"dry_run": True}
|
|
r = self.s.patch(self._url(path), json=payload, timeout=60)
|
|
self._raise(r)
|
|
return r.json()
|
|
|
|
# -- Konsumvorgaenge -----------------------------------------------------
|
|
|
|
@staticmethod
|
|
def _rows(data):
|
|
"""/api/tasks/ liefert je nach Stand eine Liste oder ein Seitenobjekt."""
|
|
if isinstance(data, list):
|
|
return data
|
|
if isinstance(data, dict):
|
|
return data.get("results") or []
|
|
return []
|
|
|
|
@staticmethod
|
|
def _doc_id(task):
|
|
"""
|
|
Die Dokument-ID steht je nach Stand an unterschiedlichen Stellen.
|
|
In 3.1.0 gleich zweimal: verschachtelt in result_data.document_id
|
|
und als Liste in related_document_ids.
|
|
"""
|
|
def as_int(val):
|
|
try:
|
|
return int(val)
|
|
except (TypeError, ValueError):
|
|
return None
|
|
|
|
# Listenfelder zuerst - dort steht die ID in 3.1.0
|
|
for key in ("related_document_ids", "document_ids"):
|
|
val = task.get(key)
|
|
if isinstance(val, (list, tuple)) and val:
|
|
got = as_int(val[0])
|
|
if got:
|
|
return got
|
|
|
|
# Verschachtelt
|
|
nested = task.get("result_data")
|
|
if isinstance(nested, dict):
|
|
for key in ("document_id", "related_document", "id"):
|
|
got = as_int(nested.get(key))
|
|
if got:
|
|
return got
|
|
|
|
# Oberste Ebene
|
|
for key in ("related_document", "document_id", "related_document_id"):
|
|
got = as_int(task.get(key))
|
|
if got:
|
|
return got
|
|
|
|
# Notnagel: ID aus einem Ergebnistext ziehen
|
|
m = re.search(r"[Dd]ocument\D+(\d+)", str(task.get("result") or ""))
|
|
return int(m.group(1)) if m else None
|
|
|
|
@staticmethod
|
|
def _error_text(task, status):
|
|
"""
|
|
Die Fehlermeldung steht je nach Stand an verschiedenen Stellen.
|
|
Ein blosses "Task FAILURE" im Journal ist wertlos - dann muss man
|
|
die Ursache spaeter muehsam aus den Container-Logs rekonstruieren.
|
|
"""
|
|
rd = task.get("result_data") or {}
|
|
for src in (task.get("result"), rd.get("error"), rd.get("exc_message"),
|
|
rd.get("exc_type"), task.get("status_display")):
|
|
if src:
|
|
txt = " ".join(str(src).split())
|
|
if txt and txt.lower() not in ("failure", "failed"):
|
|
return txt[:500]
|
|
fn = (task.get("input_data") or {}).get("filename", "?")
|
|
return f"Task {status} ohne Meldung (Datei: {fn}, task_id: {task.get('task_id')})"
|
|
|
|
def find_task(self, task_id):
|
|
"""
|
|
Sucht einen Task. Greift der Serverfilter nicht, wird die
|
|
Gesamtliste durchsucht - der Filtername hat zwischen Versionen
|
|
gewechselt, und ein stiller Fehlschlag laesst das Skript sonst
|
|
bis zum Timeout warten.
|
|
"""
|
|
try:
|
|
rows = self._rows(self.get("tasks/", task_id=task_id))
|
|
for t in rows:
|
|
if str(t.get("task_id")) == str(task_id):
|
|
return t
|
|
if len(rows) == 1 and not rows[0].get("task_id"):
|
|
return rows[0]
|
|
except Exception: # noqa: BLE001
|
|
rows = []
|
|
for t in self._rows(self.get("tasks/")):
|
|
if str(t.get("task_id")) == str(task_id):
|
|
return t
|
|
return None
|
|
|
|
def wait_for_task(self, task_id, log=None):
|
|
"""Wartet auf einen Konsumvorgang und liefert die Dokument-ID."""
|
|
if self.dry_run:
|
|
return None
|
|
if not task_id or not isinstance(task_id, str):
|
|
raise RuntimeError(f"Unerwartete Antwort auf den Upload: {task_id!r}")
|
|
|
|
deadline = time.time() + TASK_TIMEOUT
|
|
waited = 0
|
|
while time.time() < deadline:
|
|
task = self.find_task(task_id)
|
|
if task:
|
|
# Die API liefert den Status kleingeschrieben ("success"),
|
|
# aeltere Staende grossgeschrieben. Vergleich deshalb
|
|
# unabhaengig von der Schreibweise.
|
|
status = str(task.get("status") or "").upper()
|
|
if status in ("SUCCESS", "SUCCEEDED"):
|
|
doc = self._doc_id(task)
|
|
if doc:
|
|
return doc
|
|
raise RuntimeError(f"Task erfolgreich, aber ohne Dokument-ID: {task}")
|
|
if status in ("FAILURE", "FAILED", "REVOKED"):
|
|
raise RuntimeError(self._error_text(task, status))
|
|
time.sleep(TASK_POLL_INTERVAL)
|
|
waited += TASK_POLL_INTERVAL
|
|
if log and waited % 30 == 0:
|
|
state = task.get("status") if task else "nicht gefunden"
|
|
log(f" ... warte seit {waited}s (Task: {state})")
|
|
raise TimeoutError(
|
|
f"Task {task_id} nach {TASK_TIMEOUT}s ohne Ergebnis. "
|
|
f"Pruefen: curl -H \"Authorization: Token \\$PT\" "
|
|
f"'{self.base}/api/tasks/?task_id={task_id}'"
|
|
)
|
|
|
|
def post_document(self, path: Path, fields: dict):
|
|
with path.open("rb") as fh:
|
|
data = {k: v for k, v in fields.items() if v not in (None, "", [])}
|
|
files = {"document": (path.name, fh, "application/octet-stream")}
|
|
tags = data.pop("tags", [])
|
|
payload = [(k, str(v)) for k, v in data.items()]
|
|
payload += [("tags", str(t)) for t in tags]
|
|
r = self.post("documents/post_document/", files=files, data=payload)
|
|
return None if isinstance(r, dict) and r.get("dry_run") else r
|
|
|
|
def update_version(self, doc_id: int, path: Path, label: str):
|
|
with path.open("rb") as fh:
|
|
files = {"document": (path.name, fh, "application/octet-stream")}
|
|
r = self.post(
|
|
f"documents/{doc_id}/update_version/",
|
|
files=files,
|
|
data={"version_label": label},
|
|
)
|
|
return None if isinstance(r, dict) and r.get("dry_run") else r
|
|
|
|
|
|
# ------------------------------------------------------------------ Auflöser
|
|
|
|
class Resolver:
|
|
"""Legt Stammdaten bei Bedarf an und merkt sich die IDs."""
|
|
|
|
def __init__(self, api: Paperless, rollen: dict, log):
|
|
self.api, self.rollen, self.log = api, rollen, log
|
|
self.tags = self._index("tags/")
|
|
self.types = self._index("document_types/")
|
|
self.correspondents = self._index("correspondents/")
|
|
self.fields = self._index("custom_fields/")
|
|
self.users = self._index("users/", key="username")
|
|
self.groups = self._index("groups/")
|
|
|
|
def _index(self, path, key="name"):
|
|
return {row[key]: row["id"] for row in self.api.get_all(path)}
|
|
|
|
def _ensure(self, cache, path, name, extra=None):
|
|
if name in cache:
|
|
return cache[name]
|
|
if self.api.dry_run:
|
|
self.log(f" [dry-run] wuerde anlegen: {path} '{name}'")
|
|
cache[name] = -1
|
|
return -1
|
|
try:
|
|
row = self.api.post(path, json={"name": name, **(extra or {})})
|
|
except Exception as exc: # noqa: BLE001
|
|
# Paperless normalisiert Namen beim Speichern (Leerzeichen).
|
|
# Ein 400 heisst deshalb meist: existiert schon unter leicht
|
|
# anderem Namen. Liste neu einlesen und erneut nachsehen.
|
|
fresh = self._index(path)
|
|
hit = fresh.get(name) or fresh.get(name.strip())
|
|
if hit is None:
|
|
low = {k.strip().casefold(): v for k, v in fresh.items()}
|
|
hit = low.get(name.strip().casefold())
|
|
if hit is None:
|
|
raise RuntimeError(f"{path} '{name}' nicht anlegbar: {exc}") from exc
|
|
cache.clear(); cache.update(fresh); cache[name] = hit
|
|
self.log(f" vorhanden: {path} '{name}' -> {hit}")
|
|
return hit
|
|
cache[name] = row["id"]
|
|
self.log(f" angelegt: {path} '{name}' -> {row['id']}")
|
|
return row["id"]
|
|
|
|
def tag(self, name):
|
|
return self._ensure(self.tags, "tags/", name)
|
|
|
|
def doc_type(self, name):
|
|
return self._ensure(self.types, "document_types/", name)
|
|
|
|
def custom_field_typed(self, name, paperless_type):
|
|
"""Feld mit ausdruecklich angegebenem Paperless-Datentyp."""
|
|
return self._ensure(self.fields, "custom_fields/", name,
|
|
{"data_type": paperless_type})
|
|
|
|
def custom_field(self, name, ecodms_type):
|
|
mapping = {"Date": "date", "CheckBox": "boolean", "String": "string"}
|
|
return self._ensure(
|
|
self.fields,
|
|
"custom_fields/",
|
|
name,
|
|
{"data_type": mapping.get(ecodms_type, "string")},
|
|
)
|
|
|
|
# -- Rechte --------------------------------------------------------------
|
|
|
|
def permissions_for(self, roles):
|
|
"""
|
|
ecoDMS-Rollen -> Paperless owner + set_permissions.
|
|
Nicht abgebildete Rollen fuehren zum Abbruch, nicht zum stillen
|
|
Verwerfen: In Paperless heisst "keine Rechte" nicht gesperrt,
|
|
sondern unbeschraenkt.
|
|
"""
|
|
cfg = self.rollen["roles"]
|
|
rights = self.rollen["rights"]
|
|
owner, view_u, view_g, chg_u, chg_g = None, set(), set(), set(), set()
|
|
|
|
for entry in roles:
|
|
role, right = entry["role"], entry["right"]
|
|
if role not in cfg:
|
|
raise KeyError(
|
|
f"Rolle '{role}' fehlt in rollen.yml. Ergaenzen oder "
|
|
f"ausdruecklich als 'type: null' eintragen."
|
|
)
|
|
spec = cfg[role]
|
|
if spec.get("type") in (None, "null"):
|
|
continue
|
|
if right not in rights:
|
|
raise KeyError(f"doc_right '{right}' fehlt in rollen.yml")
|
|
|
|
perms = rights[right]
|
|
if spec["type"] == "user":
|
|
uid = self.users.get(spec["paperless"])
|
|
if uid is None:
|
|
raise KeyError(f"Benutzer '{spec['paperless']}' existiert nicht")
|
|
if spec.get("owner") and owner is None:
|
|
owner = uid
|
|
if "view" in perms:
|
|
view_u.add(uid)
|
|
if "change" in perms:
|
|
chg_u.add(uid)
|
|
else:
|
|
gid = self.groups.get(spec["paperless"])
|
|
if gid is None:
|
|
raise KeyError(f"Gruppe '{spec['paperless']}' existiert nicht")
|
|
if "view" in perms:
|
|
view_g.add(gid)
|
|
if "change" in perms:
|
|
chg_g.add(gid)
|
|
|
|
# no_rights_policy: owner_only
|
|
if owner is None:
|
|
owner = self.users[self.rollen["default_owner"]]
|
|
if not (view_u or view_g or chg_u or chg_g):
|
|
return owner, None # nur Eigentuemer, sonst niemand
|
|
|
|
return owner, {
|
|
"view": {"users": sorted(view_u), "groups": sorted(view_g)},
|
|
"change": {"users": sorted(chg_u), "groups": sorted(chg_g)},
|
|
}
|
|
|
|
|
|
# ------------------------------------------------------------------- Import
|
|
|
|
def iso_added(raw):
|
|
"""ecoDMS savedate -> ISO-8601 mit Zeitzone."""
|
|
if not raw:
|
|
return None
|
|
txt = str(raw).strip().replace(" ", "T")
|
|
if "." in txt: # Mikrosekunden auf 6 Stellen kuerzen
|
|
head, frac = txt.split(".", 1)
|
|
txt = f"{head}.{frac[:6]}"
|
|
try:
|
|
return datetime.fromisoformat(txt).astimezone().isoformat()
|
|
except ValueError:
|
|
return None
|
|
|
|
|
|
def folder_tags(doc_or_folders, mode):
|
|
"""
|
|
Ordnerpfade -> Tagnamen.
|
|
|
|
'segments' (Vorgabe): 'Beruf/Fortbildung/Steuerberater' wird zu drei
|
|
Tags. Erlaubt Filtern nach einzelnen Ebenen, kostet aber den
|
|
Hierarchiekontext - 'Steuerberater' allein ist mehrdeutig. Bei
|
|
tiefen Baeumen entstehen entsprechend viele Tags.
|
|
'path': ein Tag je Ordnerpfad, Hierarchie bleibt lesbar.
|
|
|
|
Bei Mehrfachklassifizierung werden die Segmente vereinigt, Reihenfolge
|
|
bleibt stabil.
|
|
"""
|
|
# Stufe 1 liefert die Ebenen als Liste (folder_parts). Ordnernamen
|
|
# koennen Schraegstriche enthalten, ein Zerlegen am Schraegstrich waere
|
|
# deshalb falsch. Die Liste "folders" ist nur fuer die Anzeige.
|
|
if isinstance(doc_or_folders, dict):
|
|
parts_list = doc_or_folders.get("folder_parts")
|
|
if parts_list is None: # aeltere Manifeste
|
|
parts_list = [f.split(" / ") for f in doc_or_folders.get("folders", [])]
|
|
else:
|
|
parts_list = [f.split(" / ") for f in doc_or_folders]
|
|
|
|
out, seen = [], set()
|
|
for parts in parts_list:
|
|
names = [" / ".join(parts)] if mode == "path" else parts
|
|
for n in names:
|
|
n = str(n).strip()
|
|
if n and n.casefold() not in seen:
|
|
seen.add(n.casefold())
|
|
out.append(n)
|
|
return out
|
|
|
|
|
|
def import_document(doc, api, res, jr, archive: Path, log, tag_mode="segments"):
|
|
docid = doc["ecodms_docid"]
|
|
state = jr.state(docid) or {}
|
|
|
|
if state.get("status") == "done":
|
|
log(f" uebersprungen (bereits importiert als {state['paperless_id']})")
|
|
return state["paperless_id"]
|
|
|
|
versions = doc["versions"]
|
|
if not versions:
|
|
raise RuntimeError("keine Version im Manifest")
|
|
|
|
jr.start(docid, doc["tranche"], len(versions))
|
|
|
|
# -- Tags: Ordnerebenen + Status
|
|
tag_ids = [res.tag(t) for t in folder_tags(doc, tag_mode)]
|
|
if doc.get("status"):
|
|
tag_ids.append(res.tag(f"Status: {doc['status']}"))
|
|
|
|
paperless_id = state.get("paperless_id")
|
|
start_at = state.get("versions_done") or 0
|
|
|
|
# -- Version 1 als Wurzeldokument
|
|
if not paperless_id:
|
|
v1 = versions[0]
|
|
f = archive / v1["file"]
|
|
if not f.exists():
|
|
raise FileNotFoundError(f)
|
|
log(f" v1: {v1['file']}")
|
|
task = api.post_document(
|
|
f,
|
|
{
|
|
"title": doc["title"][:127],
|
|
"created": doc.get("created") or None,
|
|
"document_type": res.doc_type(doc["document_type"]),
|
|
"tags": tag_ids,
|
|
},
|
|
)
|
|
paperless_id = api.wait_for_task(task, log) if task else None
|
|
jr.update(docid, paperless_id=paperless_id, versions_done=1,
|
|
status="root_created")
|
|
if paperless_id:
|
|
jr.remember_time(paperless_id, docid, 1, iso_added(v1["saved_at"]))
|
|
start_at = 1
|
|
|
|
# -- Version 2..n strikt seriell anhaengen
|
|
for v in versions[start_at:]:
|
|
f = archive / v["file"]
|
|
if not f.exists():
|
|
raise FileNotFoundError(f)
|
|
label = f"v{v['version']}"
|
|
if v.get("saved_at"):
|
|
label += f" ({str(v['saved_at'])[:10]})"
|
|
log(f" v{v['version']}: {v['file']} [{label}]")
|
|
task = api.update_version(paperless_id, f, label)
|
|
vid = api.wait_for_task(task, log) if task else None
|
|
if vid:
|
|
jr.remember_time(vid, docid, v["version"], iso_added(v["saved_at"]))
|
|
jr.update(docid, versions_done=v["version"])
|
|
|
|
# -- Custom Fields
|
|
cf = []
|
|
for name, meta in doc.get("custom_fields", {}).items():
|
|
cf.append({"field": res.custom_field(name, meta["type"]),
|
|
"value": meta["value"]})
|
|
cf.append({"field": res.custom_field("ecoDMS-ID", "String"),
|
|
"value": str(docid)})
|
|
|
|
# -- Rechte: einmal auf die Wurzel, Versionen erben ueber root_document
|
|
owner, perms = res.permissions_for(doc["permissions"])
|
|
payload = {
|
|
"owner": owner,
|
|
"custom_fields": cf,
|
|
"archive_serial_number": docid, # ecoDMS-docid als ASN
|
|
}
|
|
if perms:
|
|
payload["set_permissions"] = perms
|
|
|
|
if paperless_id:
|
|
try:
|
|
api.patch(f"documents/{paperless_id}/", payload)
|
|
except Exception as exc: # noqa: BLE001
|
|
# Die ASN ist in Paperless eindeutig. Ein Konflikt darf nicht
|
|
# das ganze Dokument kosten - lieber ohne ASN weitermachen und
|
|
# den Fall protokollieren.
|
|
if "archive_serial_number" in str(exc) or "serial" in str(exc).lower():
|
|
log(f" ASN {docid} abgelehnt ({exc}) - Import ohne ASN")
|
|
payload.pop("archive_serial_number")
|
|
api.patch(f"documents/{paperless_id}/", payload)
|
|
else:
|
|
raise
|
|
|
|
jr.update(docid, status="done", last_error=None)
|
|
return paperless_id
|
|
|
|
|
|
def run_import(args, manifest, api, res, jr, log):
|
|
archive = Path(manifest["source"])
|
|
docs = manifest["documents"]
|
|
if args.limit:
|
|
docs = docs[: args.limit]
|
|
|
|
ok = failed = 0
|
|
for i, doc in enumerate(docs, 1):
|
|
log(f"[{i}/{len(docs)}] docid {doc['ecodms_docid']}: {doc['title'][:60]}")
|
|
try:
|
|
import_document(doc, api, res, jr, archive, log, args.tag_mode)
|
|
ok += 1
|
|
except Exception as exc: # noqa: BLE001
|
|
failed += 1
|
|
jr.update(doc["ecodms_docid"], status="failed", last_error=str(exc))
|
|
log(f" FEHLER: {exc}")
|
|
if args.stop_on_error:
|
|
raise
|
|
log(f"\nFertig: {ok} erfolgreich, {failed} fehlgeschlagen.")
|
|
|
|
|
|
# --------------------------------------------------------- Zeitstempel-Skript
|
|
|
|
TIMESTAMP_TEMPLATE = '''\
|
|
# Erzeugt von paperless_import.py --emit-timestamps
|
|
#
|
|
# Ausfuehren mit:
|
|
# sudo docker compose exec -T webserver python manage.py shell < fix_added.py
|
|
#
|
|
# update() statt save(), damit auto_now-Felder und Signal-Handler nicht
|
|
# greifen - die wuerden den Wert wieder ueberschreiben und nebenbei
|
|
# Reindexierungen ausloesen.
|
|
from documents.models import Document
|
|
from django.utils.dateparse import parse_datetime
|
|
|
|
ROWS = {rows}
|
|
|
|
changed = missing = 0
|
|
for pk, added in ROWS:
|
|
n = Document.objects.filter(pk=pk).update(added=parse_datetime(added))
|
|
if n:
|
|
changed += 1
|
|
else:
|
|
missing += 1
|
|
print(f"{{changed}} Zeitstempel gesetzt, {{missing}} Dokumente nicht gefunden")
|
|
'''
|
|
|
|
|
|
def emit_timestamps(jr):
|
|
rows = [
|
|
(r["paperless_id"], r["added"])
|
|
for r in jr.con.execute(
|
|
"SELECT paperless_id, added FROM version_times "
|
|
"WHERE added IS NOT NULL ORDER BY ecodms_docid, version"
|
|
)
|
|
]
|
|
print(TIMESTAMP_TEMPLATE.format(rows=repr(rows)))
|
|
|
|
|
|
# ------------------------------------------------------------ Verknuepfungen
|
|
|
|
def resolve_links(path: Path, api, res, jr, log):
|
|
"""
|
|
Phase 2, erst nach ALLEN Tranchen: Erst jetzt sind alle docids im
|
|
Journal, also auch die aus spaeteren Tranchen.
|
|
|
|
CSV-Format, eine Zeile je Paar: src_docid,dst_docid
|
|
"""
|
|
with path.open() as fh:
|
|
for row in csv.reader(fh):
|
|
if len(row) >= 2 and row[0].strip().isdigit():
|
|
jr.add_link(int(row[0]), int(row[1]))
|
|
|
|
field_id = res.custom_field("ecoDMS-Verknuepfung", "documentlink")
|
|
pending = list(jr.con.execute("SELECT * FROM links WHERE resolved = 0"))
|
|
log(f"{len(pending)} Verknuepfungen aufzuloesen")
|
|
|
|
# Beide Richtungen sammeln - die Automatik ist beim Bulk-Edit
|
|
# nachweislich einseitig (Issue #8960).
|
|
partners: dict[int, set[int]] = {}
|
|
unresolved = 0
|
|
for row in pending:
|
|
a = jr.state(row["src_docid"])
|
|
b = jr.state(row["dst_docid"])
|
|
if not (a and b and a["paperless_id"] and b["paperless_id"]):
|
|
unresolved += 1
|
|
continue
|
|
partners.setdefault(a["paperless_id"], set()).add(b["paperless_id"])
|
|
partners.setdefault(b["paperless_id"], set()).add(a["paperless_id"])
|
|
|
|
for pid, others in partners.items():
|
|
api.patch(
|
|
f"documents/{pid}/",
|
|
{"custom_fields": [{"field": field_id, "value": sorted(others)}]},
|
|
)
|
|
jr.con.execute("UPDATE links SET resolved = 1 WHERE resolved = 0")
|
|
jr.con.commit()
|
|
log(f"{len(partners)} Dokumente verknuepft, {unresolved} ohne Gegenstueck")
|
|
|
|
|
|
# ----------------------------------------------------------------- Pruefung
|
|
|
|
def verify(api, jr, log):
|
|
"""
|
|
Prueft das Journal gegen die Instanz. Holt die Dokumente in einem
|
|
Stapelabruf statt einzeln - bei ueber tausend Dokumenten waere ein
|
|
Aufruf je Dokument unbrauchbar langsam.
|
|
"""
|
|
rows = list(jr.con.execute("SELECT * FROM mapping ORDER BY ecodms_docid"))
|
|
log(f"Journal: {len(rows)} Dokumente")
|
|
for status in ("done", "root_created", "pending", "uebersprungen", "failed"):
|
|
n = sum(1 for r in rows if r["status"] == status)
|
|
if n:
|
|
log(f" {status:14s} {n}")
|
|
|
|
log("\n lade Dokumente aus Paperless ...")
|
|
live = {}
|
|
for d in api.get_all("documents/",
|
|
fields="id,owner,archive_serial_number,versions"):
|
|
live[d["id"]] = d
|
|
log(f" {len(live)} Dokumente in der Instanz\n")
|
|
|
|
problems = 0
|
|
def flag(msg):
|
|
nonlocal problems
|
|
problems += 1
|
|
if problems <= 40:
|
|
log(f" {msg}")
|
|
|
|
seen = set()
|
|
for r in rows:
|
|
if r["status"] not in ("done", "root_created"):
|
|
continue
|
|
pid = r["paperless_id"]
|
|
if not pid:
|
|
flag(f"docid {r['ecodms_docid']}: ohne Paperless-ID")
|
|
continue
|
|
d = live.get(pid)
|
|
if d is None:
|
|
flag(f"docid {r['ecodms_docid']}: Dokument {pid} nicht in der Instanz")
|
|
continue
|
|
seen.add(pid)
|
|
n = len(d.get("versions") or []) or 1
|
|
if n != r["versions_total"]:
|
|
flag(f"docid {r['ecodms_docid']}: {n} Versionen, "
|
|
f"{r['versions_total']} erwartet")
|
|
if d.get("owner") is None:
|
|
flag(f"docid {r['ecodms_docid']}: KEIN EIGENTUEMER "
|
|
f"(waere unbeschraenkt sichtbar)")
|
|
asn = d.get("archive_serial_number")
|
|
if asn is not None and int(asn) != r["ecodms_docid"]:
|
|
flag(f"docid {r['ecodms_docid']}: ASN {asn} weicht ab")
|
|
|
|
# Dokumente in der Instanz, die keine Wurzel aus dem Journal sind und
|
|
# auch keine Version davon - Kandidaten fuer doppelte Anlage.
|
|
versions_of_known = {
|
|
v["id"] for pid in seen for v in (live.get(pid, {}).get("versions") or [])
|
|
}
|
|
orphans = sorted(set(live) - seen - versions_of_known)
|
|
if orphans:
|
|
log(f"\n {len(orphans)} Dokumente ohne Journaleintrag "
|
|
f"(moegliche Doppelanlage): {orphans[:20]}")
|
|
problems += len(orphans)
|
|
|
|
if problems > 40:
|
|
log(f" ... und {problems - 40} weitere")
|
|
log(f"\n{problems} Auffaelligkeiten." if problems else "\nKeine Auffaelligkeiten.")
|
|
|
|
|
|
# --------------------------------------------------------------------- Main
|
|
|
|
def main():
|
|
ap = argparse.ArgumentParser(description=__doc__,
|
|
formatter_class=argparse.RawDescriptionHelpFormatter)
|
|
ap.add_argument("manifest", nargs="?", type=Path, help="JSON aus Stufe 1")
|
|
ap.add_argument("--api", default="http://localhost:8000")
|
|
ap.add_argument("--token", help="oder Umgebungsvariable PT")
|
|
ap.add_argument("--rollen", type=Path, default=Path("rollen.yml"))
|
|
ap.add_argument("--journal", type=Path, default=Path("migration.sqlite"))
|
|
ap.add_argument("--dry-run", action="store_true")
|
|
ap.add_argument("--limit", type=int, help="nur die ersten N Dokumente")
|
|
ap.add_argument("--stop-on-error", action="store_true")
|
|
ap.add_argument("--tag-mode", choices=("segments", "path"), default="segments",
|
|
help="segments: jede Ordnerebene ein eigener Tag (Vorgabe). "
|
|
"path: ein Tag je vollstaendigem Ordnerpfad.")
|
|
ap.add_argument("--setup", action="store_true",
|
|
help="nur Stammdaten anlegen, nichts importieren")
|
|
ap.add_argument("--emit-timestamps", action="store_true")
|
|
ap.add_argument("--links", type=Path, metavar="CSV")
|
|
ap.add_argument("--verify", action="store_true")
|
|
args = ap.parse_args()
|
|
|
|
log = print
|
|
jr = Journal(args.journal)
|
|
|
|
if args.emit_timestamps:
|
|
emit_timestamps(jr)
|
|
return
|
|
|
|
import os
|
|
token = args.token or os.environ.get("PT")
|
|
if not token:
|
|
sys.exit("Token fehlt: --token oder export PT=...")
|
|
|
|
api = Paperless(args.api, token, dry_run=args.dry_run)
|
|
rollen = yaml.safe_load(args.rollen.read_text(encoding="utf-8"))
|
|
res = Resolver(api, rollen, log)
|
|
|
|
if args.verify:
|
|
verify(api, jr, log)
|
|
return
|
|
|
|
if args.links:
|
|
resolve_links(args.links, api, res, jr, log)
|
|
return
|
|
|
|
if not args.manifest:
|
|
sys.exit("Manifest fehlt.")
|
|
manifest = json.loads(args.manifest.read_text(encoding="utf-8"))
|
|
|
|
log(f"Tranche {manifest['tranche']}: "
|
|
f"{manifest['counts']['documents']} Dokumente, "
|
|
f"{manifest['counts']['versions']} Versionen")
|
|
if args.dry_run:
|
|
log("TROCKENLAUF - es wird nichts geschrieben\n")
|
|
|
|
# Rollen vorab pruefen, bevor die erste Datei hochgeht
|
|
seen = {p["role"] for d in manifest["documents"] for p in d["permissions"]}
|
|
unknown = seen - set(rollen["roles"])
|
|
if unknown:
|
|
sys.exit(f"Unbekannte Rollen in rollen.yml ergaenzen: {sorted(unknown)}")
|
|
|
|
if args.setup:
|
|
for d in manifest["documents"]:
|
|
for t in folder_tags(d, args.tag_mode):
|
|
res.tag(t)
|
|
if d.get("status"):
|
|
res.tag(f"Status: {d['status']}")
|
|
res.doc_type(d["document_type"])
|
|
for name, meta in d.get("custom_fields", {}).items():
|
|
res.custom_field(name, meta["type"])
|
|
# Felder aus rollen.yml unabhaengig vom Tranchen-Inhalt anlegen,
|
|
# damit Namen und Typen nicht vom Zufall der ersten Tranche abhaengen.
|
|
for name, dtype in (rollen.get("custom_fields") or {}).items():
|
|
res.custom_field_typed(name, dtype)
|
|
log("Stammdaten angelegt.")
|
|
return
|
|
|
|
run_import(args, manifest, api, res, jr, log)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|