Initial commit
This commit is contained in:
863
paperless_import.py
Normal file
863
paperless_import.py
Normal file
@@ -0,0 +1,863 @@
|
||||
#!/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()
|
||||
Reference in New Issue
Block a user