255 lines
13 KiB
Python
255 lines
13 KiB
Python
|
|
"""Portable, checksummed Project CT archive bundles. Never contacts vTiger."""
|
||
|
|
|
||
|
|
from __future__ import annotations
|
||
|
|
|
||
|
|
import hashlib
|
||
|
|
import json
|
||
|
|
import os
|
||
|
|
import socket
|
||
|
|
import tempfile
|
||
|
|
import zipfile
|
||
|
|
from datetime import datetime, timezone
|
||
|
|
from typing import Any, Iterable
|
||
|
|
|
||
|
|
from psycopg2.extras import RealDictCursor
|
||
|
|
|
||
|
|
from app.core.database import get_db_connection, release_db_connection
|
||
|
|
|
||
|
|
BUNDLE_SCHEMA_VERSION = 1
|
||
|
|
ENTRY_NAMES = ("versions.jsonl", "records.jsonl", "relations.jsonl", "checkpoints.jsonl", "files.jsonl")
|
||
|
|
|
||
|
|
|
||
|
|
def file_sha256(path: str) -> str:
|
||
|
|
digest = hashlib.sha256()
|
||
|
|
with open(path, "rb") as source:
|
||
|
|
for chunk in iter(lambda: source.read(1024 * 1024), b""):
|
||
|
|
digest.update(chunk)
|
||
|
|
return digest.hexdigest()
|
||
|
|
|
||
|
|
|
||
|
|
def _json_line(value: Any) -> bytes:
|
||
|
|
return (json.dumps(value, ensure_ascii=False, sort_keys=True, default=str) + "\n").encode("utf-8")
|
||
|
|
|
||
|
|
|
||
|
|
def _write_cursor_entry(archive: zipfile.ZipFile, name: str, cursor, batch_size: int = 500) -> tuple[int, str]:
|
||
|
|
digest = hashlib.sha256()
|
||
|
|
count = 0
|
||
|
|
with archive.open(name, "w") as target:
|
||
|
|
while True:
|
||
|
|
rows = cursor.fetchmany(batch_size)
|
||
|
|
if not rows:
|
||
|
|
break
|
||
|
|
for row in rows:
|
||
|
|
raw = _json_line(dict(row))
|
||
|
|
target.write(raw)
|
||
|
|
digest.update(raw)
|
||
|
|
count += 1
|
||
|
|
return count, digest.hexdigest()
|
||
|
|
|
||
|
|
|
||
|
|
def create_archive_bundle(through_version_id: int, initiated_by: int | None = None) -> tuple[str, dict]:
|
||
|
|
fd, path = tempfile.mkstemp(prefix="project-ct-", suffix=".zip")
|
||
|
|
os.close(fd)
|
||
|
|
conn = get_db_connection()
|
||
|
|
manifest = {
|
||
|
|
"schema_version": BUNDLE_SCHEMA_VERSION,
|
||
|
|
"created_at": datetime.now(timezone.utc).isoformat(),
|
||
|
|
"source_instance": socket.gethostname(),
|
||
|
|
"through_version_id": through_version_id,
|
||
|
|
"entries": {},
|
||
|
|
"files": {},
|
||
|
|
"contains_files": True,
|
||
|
|
"document_policy": "metadata_and_original_bytes",
|
||
|
|
}
|
||
|
|
try:
|
||
|
|
with zipfile.ZipFile(path, "w", compression=zipfile.ZIP_DEFLATED, compresslevel=6) as archive:
|
||
|
|
queries = {
|
||
|
|
"versions.jsonl": ("SELECT * FROM vtiger_archive_versions WHERE id<=%s ORDER BY id", (through_version_id,)),
|
||
|
|
"records.jsonl": ("SELECT * FROM vtiger_archive_records WHERE version_id<=%s ORDER BY module,vtiger_id,revision_no", (through_version_id,)),
|
||
|
|
"relations.jsonl": ("SELECT * FROM vtiger_archive_relations WHERE version_id<=%s ORDER BY id", (through_version_id,)),
|
||
|
|
"checkpoints.jsonl": (
|
||
|
|
"""SELECT module,MAX(source_modified_at) AS last_modified_at,MAX(vtiger_id) AS last_vtiger_id,
|
||
|
|
%s::bigint AS last_successful_version_id,NOW() AS updated_at
|
||
|
|
FROM vtiger_archive_records WHERE version_id<=%s GROUP BY module ORDER BY module""",
|
||
|
|
(through_version_id, through_version_id),
|
||
|
|
),
|
||
|
|
"files.jsonl": (
|
||
|
|
"""SELECT id,version_id,source_module,document_vtiger_id,resource_vtiger_id,filename,content_type,
|
||
|
|
size_bytes,content_sha256,archived_at
|
||
|
|
FROM vtiger_archive_files WHERE version_id<=%s ORDER BY id""",
|
||
|
|
(through_version_id,),
|
||
|
|
),
|
||
|
|
}
|
||
|
|
for index, (name, (sql, params)) in enumerate(queries.items()):
|
||
|
|
cursor = conn.cursor(name=f"project_ct_export_{index}", cursor_factory=RealDictCursor)
|
||
|
|
cursor.itersize = 500
|
||
|
|
cursor.execute(sql, params)
|
||
|
|
count, digest = _write_cursor_entry(archive, name, cursor)
|
||
|
|
cursor.close()
|
||
|
|
manifest["entries"][name] = {"count": count, "sha256": digest}
|
||
|
|
file_rows = conn.cursor(cursor_factory=RealDictCursor)
|
||
|
|
file_rows.execute(
|
||
|
|
"""SELECT DISTINCT ON (content_sha256) content_sha256,content,size_bytes
|
||
|
|
FROM vtiger_archive_files WHERE version_id<=%s ORDER BY content_sha256,id""",
|
||
|
|
(through_version_id,),
|
||
|
|
)
|
||
|
|
for row in file_rows:
|
||
|
|
content = bytes(row["content"])
|
||
|
|
digest = hashlib.sha256(content).hexdigest()
|
||
|
|
if digest != row["content_sha256"]:
|
||
|
|
raise ValueError(f"Lokal filchecksum er ugyldig: {row['content_sha256']}")
|
||
|
|
entry_name = f"files/{digest}"
|
||
|
|
archive.writestr(entry_name, content)
|
||
|
|
manifest["files"][digest] = {"entry": entry_name, "size": len(content), "sha256": digest}
|
||
|
|
file_rows.close()
|
||
|
|
archive.writestr("manifest.json", json.dumps(manifest, ensure_ascii=False, sort_keys=True, indent=2, default=str))
|
||
|
|
manifest["bundle_sha256"] = file_sha256(path)
|
||
|
|
return path, manifest
|
||
|
|
except Exception:
|
||
|
|
if os.path.exists(path):
|
||
|
|
os.unlink(path)
|
||
|
|
raise
|
||
|
|
finally:
|
||
|
|
release_db_connection(conn)
|
||
|
|
|
||
|
|
|
||
|
|
def _verify_entry(archive: zipfile.ZipFile, manifest: dict, name: str) -> None:
|
||
|
|
expected = ((manifest.get("entries") or {}).get(name) or {}).get("sha256")
|
||
|
|
if not expected:
|
||
|
|
raise ValueError(f"Manifest mangler checksum for {name}")
|
||
|
|
digest = hashlib.sha256()
|
||
|
|
with archive.open(name) as source:
|
||
|
|
for chunk in iter(lambda: source.read(1024 * 1024), b""):
|
||
|
|
digest.update(chunk)
|
||
|
|
if digest.hexdigest() != expected:
|
||
|
|
raise ValueError(f"Checksumfejl i {name}")
|
||
|
|
|
||
|
|
|
||
|
|
def _rows(archive: zipfile.ZipFile, name: str) -> Iterable[dict]:
|
||
|
|
with archive.open(name) as source:
|
||
|
|
for line_number, raw in enumerate(source, 1):
|
||
|
|
if not raw.strip():
|
||
|
|
continue
|
||
|
|
try:
|
||
|
|
yield json.loads(raw)
|
||
|
|
except json.JSONDecodeError as exc:
|
||
|
|
raise ValueError(f"Ugyldig JSON i {name}, linje {line_number}") from exc
|
||
|
|
|
||
|
|
|
||
|
|
def import_archive_bundle(path: str, initiated_by: int | None = None) -> dict:
|
||
|
|
bundle_digest = file_sha256(path)
|
||
|
|
conn = get_db_connection()
|
||
|
|
counts = {"versions": 0, "records": 0, "records_skipped": 0, "relations": 0, "checkpoints": 0}
|
||
|
|
try:
|
||
|
|
with zipfile.ZipFile(path) as archive:
|
||
|
|
try:
|
||
|
|
manifest = json.loads(archive.read("manifest.json"))
|
||
|
|
except (KeyError, json.JSONDecodeError) as exc:
|
||
|
|
raise ValueError("Filen er ikke en gyldig Project CT-arkivpakke") from exc
|
||
|
|
if manifest.get("schema_version") != BUNDLE_SCHEMA_VERSION:
|
||
|
|
raise ValueError(f"Ikke-understøttet arkivformat: {manifest.get('schema_version')}")
|
||
|
|
if manifest.get("contains_files") is not True:
|
||
|
|
raise ValueError("Pakken mangler originale vTiger-dokumentfiler")
|
||
|
|
for name in ENTRY_NAMES:
|
||
|
|
_verify_entry(archive, manifest, name)
|
||
|
|
for digest, info in (manifest.get("files") or {}).items():
|
||
|
|
entry_name = info.get("entry")
|
||
|
|
if not entry_name or entry_name not in archive.namelist():
|
||
|
|
raise ValueError(f"Pakken mangler dokumentfil {digest}")
|
||
|
|
actual = hashlib.sha256(archive.read(entry_name)).hexdigest()
|
||
|
|
if actual != digest or actual != info.get("sha256"):
|
||
|
|
raise ValueError(f"Checksumfejl i dokumentfil {digest}")
|
||
|
|
|
||
|
|
with conn.cursor(cursor_factory=RealDictCursor) as cursor:
|
||
|
|
version_map: dict[int, int] = {}
|
||
|
|
for row in _rows(archive, "versions.jsonl"):
|
||
|
|
cursor.execute(
|
||
|
|
"""INSERT INTO vtiger_archive_versions
|
||
|
|
(sync_kind,status,started_at,completed_at,source_cutoff,module_counts,warnings,
|
||
|
|
critical_errors,control_report,control_approved_at,raw_export_sha256)
|
||
|
|
VALUES (%s,%s,%s,%s,%s,%s::jsonb,%s::jsonb,%s::jsonb,%s::jsonb,%s,%s) RETURNING id""",
|
||
|
|
(row["sync_kind"], row["status"], row["started_at"], row.get("completed_at"), row["source_cutoff"],
|
||
|
|
json.dumps(row.get("module_counts") or {}), json.dumps(row.get("warnings") or []),
|
||
|
|
json.dumps(row.get("critical_errors") or []), json.dumps(row.get("control_report") or {}),
|
||
|
|
row.get("control_approved_at"), row.get("raw_export_sha256")),
|
||
|
|
)
|
||
|
|
version_map[int(row["id"])] = int(cursor.fetchone()["id"])
|
||
|
|
counts["versions"] += 1
|
||
|
|
|
||
|
|
for row in _rows(archive, "records.jsonl"):
|
||
|
|
cursor.execute(
|
||
|
|
"SELECT id FROM vtiger_archive_records WHERE module=%s AND vtiger_id=%s AND payload_sha256=%s",
|
||
|
|
(row["module"], row["vtiger_id"], row["payload_sha256"]),
|
||
|
|
)
|
||
|
|
if cursor.fetchone():
|
||
|
|
counts["records_skipped"] += 1
|
||
|
|
continue
|
||
|
|
cursor.execute(
|
||
|
|
"SELECT COALESCE(MAX(revision_no),0)+1 AS revision FROM vtiger_archive_records WHERE module=%s AND vtiger_id=%s",
|
||
|
|
(row["module"], row["vtiger_id"]),
|
||
|
|
)
|
||
|
|
revision = int(cursor.fetchone()["revision"])
|
||
|
|
cursor.execute(
|
||
|
|
"""INSERT INTO vtiger_archive_records
|
||
|
|
(version_id,module,vtiger_id,revision_no,source_created_at,source_modified_at,is_deleted,payload,payload_sha256,archived_at)
|
||
|
|
VALUES (%s,%s,%s,%s,%s,%s,%s,%s::jsonb,%s,%s)""",
|
||
|
|
(version_map[int(row["version_id"])], row["module"], row["vtiger_id"], revision,
|
||
|
|
row.get("source_created_at"), row.get("source_modified_at"), row.get("is_deleted", False),
|
||
|
|
json.dumps(row["payload"], ensure_ascii=False), row["payload_sha256"], row.get("archived_at")),
|
||
|
|
)
|
||
|
|
counts["records"] += 1
|
||
|
|
|
||
|
|
counts["files"] = 0
|
||
|
|
counts["files_skipped"] = 0
|
||
|
|
for row in _rows(archive, "files.jsonl"):
|
||
|
|
cursor.execute(
|
||
|
|
"""SELECT id FROM vtiger_archive_files
|
||
|
|
WHERE source_module=%s AND document_vtiger_id=%s AND resource_vtiger_id=%s AND content_sha256=%s""",
|
||
|
|
(row.get("source_module") or "Documents", row["document_vtiger_id"], row["resource_vtiger_id"], row["content_sha256"]),
|
||
|
|
)
|
||
|
|
if cursor.fetchone():
|
||
|
|
counts["files_skipped"] += 1
|
||
|
|
continue
|
||
|
|
info = manifest["files"][row["content_sha256"]]
|
||
|
|
content = archive.read(info["entry"])
|
||
|
|
cursor.execute(
|
||
|
|
"""INSERT INTO vtiger_archive_files
|
||
|
|
(version_id,source_module,document_vtiger_id,resource_vtiger_id,filename,content_type,size_bytes,
|
||
|
|
content_sha256,content,archived_at)
|
||
|
|
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
|
||
|
|
(version_map[int(row["version_id"])], row.get("source_module") or "Documents", row["document_vtiger_id"], row["resource_vtiger_id"],
|
||
|
|
row["filename"], row.get("content_type"), row["size_bytes"], row["content_sha256"], content,
|
||
|
|
row.get("archived_at")),
|
||
|
|
)
|
||
|
|
counts["files"] += 1
|
||
|
|
|
||
|
|
for row in _rows(archive, "relations.jsonl"):
|
||
|
|
cursor.execute(
|
||
|
|
"""INSERT INTO vtiger_archive_relations
|
||
|
|
(version_id,source_module,source_vtiger_id,field_name,target_vtiger_id,target_module)
|
||
|
|
VALUES (%s,%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING""",
|
||
|
|
(version_map[int(row["version_id"])], row["source_module"], row["source_vtiger_id"],
|
||
|
|
row["field_name"], row["target_vtiger_id"], row.get("target_module")),
|
||
|
|
)
|
||
|
|
counts["relations"] += cursor.rowcount
|
||
|
|
|
||
|
|
for row in _rows(archive, "checkpoints.jsonl"):
|
||
|
|
mapped_version = version_map.get(int(row["last_successful_version_id"])) if row.get("last_successful_version_id") else None
|
||
|
|
cursor.execute(
|
||
|
|
"""INSERT INTO vtiger_archive_checkpoints
|
||
|
|
(module,last_modified_at,last_vtiger_id,last_successful_version_id,updated_at)
|
||
|
|
VALUES (%s,%s,%s,%s,%s) ON CONFLICT(module) DO UPDATE SET
|
||
|
|
last_modified_at=GREATEST(vtiger_archive_checkpoints.last_modified_at,EXCLUDED.last_modified_at),
|
||
|
|
last_vtiger_id=EXCLUDED.last_vtiger_id,last_successful_version_id=EXCLUDED.last_successful_version_id,
|
||
|
|
updated_at=NOW()""",
|
||
|
|
(row["module"], row.get("last_modified_at"), row.get("last_vtiger_id"), mapped_version, row.get("updated_at")),
|
||
|
|
)
|
||
|
|
counts["checkpoints"] += 1
|
||
|
|
conn.commit()
|
||
|
|
return {"bundle_sha256": bundle_digest, "manifest": manifest, "counts": counts, "contacted_vtiger": False}
|
||
|
|
except Exception:
|
||
|
|
conn.rollback()
|
||
|
|
raise
|
||
|
|
finally:
|
||
|
|
release_db_connection(conn)
|