From d1f5eb03350ed91294c5a650c6d9b2474eeff616 Mon Sep 17 00:00:00 2001 From: Trenton H <797416+stumpylog@users.noreply.github.com> Date: Fri, 7 Aug 2026 11:50:22 -0700 Subject: [PATCH] Performance: reduce memory and I/O overhead of the document exporter during zip exports (#13490) --- src/documents/export/__init__.py | 0 src/documents/export/sinks.py | 346 +++++++++++++++ .../management/commands/document_exporter.py | 404 +++++------------- src/documents/tests/export/__init__.py | 0 src/documents/tests/export/test_sinks.py | 327 ++++++++++++++ .../tests/test_management_exporter.py | 28 +- 6 files changed, 800 insertions(+), 305 deletions(-) create mode 100644 src/documents/export/__init__.py create mode 100644 src/documents/export/sinks.py create mode 100644 src/documents/tests/export/__init__.py create mode 100644 src/documents/tests/export/test_sinks.py diff --git a/src/documents/export/__init__.py b/src/documents/export/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/src/documents/export/sinks.py b/src/documents/export/sinks.py new file mode 100644 index 000000000..43b666cbe --- /dev/null +++ b/src/documents/export/sinks.py @@ -0,0 +1,346 @@ +from __future__ import annotations + +import abc +import hashlib +import json +import os +import shutil +import tempfile +import zipfile +from contextlib import AbstractContextManager +from contextlib import contextmanager +from pathlib import Path +from pathlib import PurePosixPath +from typing import TYPE_CHECKING + +from django.conf import settings +from django.core.serializers.json import DjangoJSONEncoder + +from documents.file_handling import delete_empty_directories +from documents.utils import compute_checksum +from documents.utils import copy_file_with_basic_stats + +if TYPE_CHECKING: + from collections.abc import Iterator + from typing import TextIO + + +def _dumps(content: list | dict) -> str: + """Serialize export JSON consistently across all sinks.""" + return json.dumps(content, cls=DjangoJSONEncoder, indent=2, ensure_ascii=False) + + +class StreamingManifestWriter: + """Incrementally writes a JSON array to a text handle, one record at a time. + + Knows nothing about folders or zips: it writes the array framing and records + to whatever handle the sink's ``stream()`` yields. The sink owns the handle's + lifecycle (atomic rename, compare, spooling). + """ + + def __init__(self, handle: TextIO) -> None: + self._file = handle + self._first = True + self._file.write("[") + + def write_record(self, record: dict) -> None: + if not self._first: + self._file.write(",\n") + else: + self._first = False + self._file.write(_dumps(record)) + + def write_batch(self, records: list[dict]) -> None: + for record in records: + self.write_record(record) + + def close(self) -> None: + """Write the closing bracket. Does NOT close the handle (the sink owns it).""" + self._file.write("\n]") + + +class ExportSink(AbstractContextManager, abc.ABC): + """Destination for a document export. + + The command declares export contents via three verbs; the sink decides how to + persist each. ``arcname`` is always a relative POSIX path + (e.g. ``"manifest.json"``, ``"originals/foo.pdf"``). + + Contract: + * At most one ``stream()`` open at a time (it is the manifest); + ``add_file``/``add_json`` may be called while it is open. + * Context-manager: normal exit finalizes, an exception aborts. No partial or + failed run leaves a complete-looking artifact. + """ + + @abc.abstractmethod + def add_file( + self, + source: Path, + arcname: str, + *, + checksum: str | None = None, + ) -> None: ... + + @abc.abstractmethod + def add_json(self, content: list | dict, arcname: str) -> None: ... + + @abc.abstractmethod + def stream(self, arcname: str) -> AbstractContextManager[TextIO]: ... + + def _open(self) -> None: + """Hook called on context entry. Override as needed.""" + + @abc.abstractmethod + def _finalize(self) -> None: + """Commit on clean exit.""" + + @abc.abstractmethod + def _abort(self) -> None: + """Roll back on exception.""" + + def __enter__(self) -> ExportSink: + self._open() + return self + + def __exit__(self, exc_type, exc_val, exc_tb) -> None: + if exc_type is not None: + self._abort() + else: + self._finalize() + + +class DirectoryExportSink(ExportSink): + """Writes loose files into a target directory, with incremental sync. + + Owns the snapshot/skip/compare/prune machinery that used to live in the + command (``files_in_export_dir``, ``check_and_copy``, ``check_and_write_json``, + and the ``--delete`` pass). + """ + + def __init__( + self, + target: Path, + *, + compare_checksums: bool, + compare_json: bool, + delete: bool, + ) -> None: + self._target = target.resolve() + self._compare_checksums = compare_checksums + self._compare_json = compare_json + self._delete = delete + self._snapshot: set[Path] = set() + self._stream_open = False + + def _open(self) -> None: + for x in self._target.glob("**/*"): + if x.is_file(): + self._snapshot.add(x.resolve()) + + def add_file( + self, + source: Path, + arcname: str, + *, + checksum: str | None = None, + ) -> None: + target = (self._target / arcname).resolve() + self._snapshot.discard(target) + perform_copy = False + if target.exists(): + source_stat = source.stat() + target_stat = target.stat() + if self._compare_checksums and checksum: + perform_copy = compute_checksum(target) != checksum + elif ( + source_stat.st_mtime != target_stat.st_mtime + or source_stat.st_size != target_stat.st_size + ): + perform_copy = True + else: + perform_copy = True + if perform_copy: + target.parent.mkdir(parents=True, exist_ok=True) + copy_file_with_basic_stats(source, target) + + @staticmethod + def _content_unchanged(target: Path, new_bytes: bytes) -> bool: + """True if ``target`` already holds byte-identical content (BLAKE2b).""" + return ( + hashlib.blake2b(target.read_bytes()).hexdigest() + == hashlib.blake2b(new_bytes).hexdigest() + ) + + def add_json(self, content: list | dict, arcname: str) -> None: + target = (self._target / arcname).resolve() + json_str = _dumps(content) + perform_write = True + if target in self._snapshot: + self._snapshot.discard(target) + if self._compare_json and self._content_unchanged( + target, + json_str.encode("utf-8"), + ): + perform_write = False + if perform_write: + target.parent.mkdir(parents=True, exist_ok=True) + target.write_text(json_str, encoding="utf-8") + + @contextmanager + def stream(self, arcname: str) -> Iterator[TextIO]: + if self._stream_open: + raise RuntimeError("A stream is already open on this sink") + target = (self._target / arcname).resolve() + tmp = target.with_suffix(target.suffix + ".tmp") + target.parent.mkdir(parents=True, exist_ok=True) + handle = tmp.open("w", encoding="utf-8") + self._stream_open = True + try: + yield handle + except BaseException: + handle.close() + tmp.unlink(missing_ok=True) + raise + else: + handle.close() + self._commit_streamed_file(target, tmp) + finally: + self._stream_open = False + + def _commit_streamed_file(self, target: Path, tmp: Path) -> None: + if target in self._snapshot: + self._snapshot.discard(target) + if self._compare_json and self._content_unchanged( + target, + tmp.read_bytes(), + ): + tmp.unlink() + return + tmp.rename(target) + + def _finalize(self) -> None: + if self._delete: + for f in self._snapshot: + if not f.is_relative_to(self._target): # pragma: no cover + # Defense in depth: a symlink inside the export dir can + # resolve outside of it; never delete outside the target. + continue + f.unlink() + delete_empty_directories(f.parent, self._target) + + def _abort(self) -> None: + # Folder mode is in-place/incremental: streamed .tmp files are already + # cleaned in stream(); leave everything else intact and skip the prune. + return None + + +class ZipExportSink(ExportSink): + """Writes a single zip archive, produced atomically only on success. + + Builds into ``/.zip.tmp`` and renames to ``.zip`` on clean + finalize. The manifest stream is spooled to a temp file in SCRATCH_DIR and + added as an entry at finalize (a zip entry cannot be interleaved with others). + """ + + def __init__(self, target: Path, zip_name: str, *, delete: bool = False) -> None: + self._target = target.resolve() + self._zip_path = (self._target / zip_name).with_suffix(".zip") + self._tmp_path = self._zip_path.with_name(self._zip_path.name + ".tmp") + self._delete = delete + self._zip: zipfile.ZipFile | None = None + self._dirs: set[str] = set() + self._pending_manifest: tuple[Path, str] | None = None + self._stream_open = False + + def _open(self) -> None: + settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True) + self._zip = zipfile.ZipFile( + self._tmp_path, + "w", + compression=zipfile.ZIP_DEFLATED, + allowZip64=True, + ) + + def _ensure_dirs(self, arcname: str) -> None: + assert self._zip is not None + dir_arc = "" + for part in PurePosixPath(arcname).parts[:-1]: + dir_arc += f"{part}/" + if dir_arc not in self._dirs: + self._dirs.add(dir_arc) + self._zip.mkdir(dir_arc) + + def add_file( + self, + source: Path, + arcname: str, + *, + checksum: str | None = None, + ) -> None: + assert self._zip is not None + self._ensure_dirs(arcname) + self._zip.write(source, arcname=arcname) + + def add_json(self, content: list | dict, arcname: str) -> None: + assert self._zip is not None + self._ensure_dirs(arcname) + self._zip.writestr(arcname, _dumps(content)) + + @contextmanager + def stream(self, arcname: str) -> Iterator[TextIO]: + if self._stream_open: + raise RuntimeError("A stream is already open on this sink") + settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True) + fd, tmp_name = tempfile.mkstemp( + dir=settings.SCRATCH_DIR, + prefix="export-manifest-", + suffix=".json", + ) + tmp = Path(tmp_name) + handle = os.fdopen(fd, "w", encoding="utf-8") + self._stream_open = True + try: + yield handle + except BaseException: + handle.close() + tmp.unlink(missing_ok=True) + raise + else: + handle.close() + self._pending_manifest = (tmp, arcname) + finally: + self._stream_open = False + + def _finalize(self) -> None: + assert self._zip is not None + if self._pending_manifest is not None: + tmp, arcname = self._pending_manifest + self._ensure_dirs(arcname) + self._zip.write(tmp, arcname=arcname) + tmp.unlink(missing_ok=True) + self._pending_manifest = None + self._zip.close() + self._zip = None + if self._delete: + self._wipe_destination() + self._tmp_path.replace(self._zip_path) + + def _wipe_destination(self) -> None: + skip = {self._zip_path.resolve(), self._tmp_path.resolve()} + for item in self._target.glob("*"): + if item.resolve() in skip: + continue + if item.is_dir(): + shutil.rmtree(item) + else: + item.unlink() + + def _abort(self) -> None: + if self._zip is not None: + self._zip.close() + self._zip = None + self._tmp_path.unlink(missing_ok=True) + if self._pending_manifest is not None: + self._pending_manifest[0].unlink(missing_ok=True) + self._pending_manifest = None diff --git a/src/documents/management/commands/document_exporter.py b/src/documents/management/commands/document_exporter.py index b89faf7d9..da98e88cd 100644 --- a/src/documents/management/commands/document_exporter.py +++ b/src/documents/management/commands/document_exporter.py @@ -1,8 +1,4 @@ -import hashlib -import json import os -import shutil -import tempfile from itertools import islice from pathlib import Path from typing import TYPE_CHECKING @@ -19,7 +15,6 @@ from django.contrib.auth.models import User from django.contrib.contenttypes.models import ContentType from django.core import serializers from django.core.management.base import CommandError -from django.core.serializers.json import DjangoJSONEncoder from django.db import transaction from django.utils import timezone from filelock import FileLock @@ -34,7 +29,10 @@ if TYPE_CHECKING: if settings.AUDIT_LOG_ENABLED: from auditlog.models import LogEntry -from documents.file_handling import delete_empty_directories +from documents.export.sinks import DirectoryExportSink +from documents.export.sinks import ExportSink +from documents.export.sinks import StreamingManifestWriter +from documents.export.sinks import ZipExportSink from documents.file_handling import generate_filename from documents.management.commands.base import PaperlessCommand from documents.management.commands.mixins import CryptMixin @@ -60,8 +58,7 @@ from documents.settings import EXPORTER_ARCHIVE_NAME from documents.settings import EXPORTER_FILE_NAME from documents.settings import EXPORTER_SHARE_LINK_BUNDLE_NAME from documents.settings import EXPORTER_THUMBNAIL_NAME -from documents.utils import compute_checksum -from documents.utils import copy_file_with_basic_stats +from documents.utils import QuerySetStream from paperless import version from paperless.models import ApplicationConfiguration from paperless_mail.models import MailAccount @@ -84,87 +81,6 @@ def serialize_queryset_batched( yield serializers.serialize("python", chunk) -class StreamingManifestWriter: - """Incrementally writes a JSON array to a file, one record at a time. - - Writes to .tmp first; on close(), optionally BLAKE2b-compares - with the existing file (--compare-json) and renames or discards accordingly. - On exception, discard() deletes the tmp file and leaves the original intact. - """ - - def __init__( - self, - path: Path, - *, - compare_json: bool = False, - files_in_export_dir: "set[Path] | None" = None, - ) -> None: - self._path = path.resolve() - self._tmp_path = self._path.with_suffix(self._path.suffix + ".tmp") - self._compare_json = compare_json - self._files_in_export_dir: set[Path] = ( - files_in_export_dir if files_in_export_dir is not None else set() - ) - self._file = None - self._first = True - - def open(self) -> None: - self._path.parent.mkdir(parents=True, exist_ok=True) - self._file = self._tmp_path.open("w", encoding="utf-8") - self._file.write("[") - self._first = True - - def write_record(self, record: dict) -> None: - if not self._first: - self._file.write(",\n") - else: - self._first = False - self._file.write( - json.dumps(record, cls=DjangoJSONEncoder, indent=2, ensure_ascii=False), - ) - - def write_batch(self, records: list[dict]) -> None: - for record in records: - self.write_record(record) - - def close(self) -> None: - if self._file is None: - return - self._file.write("\n]") - self._file.close() - self._file = None - self._finalize() - - def discard(self) -> None: - if self._file is not None: - self._file.close() - self._file = None - if self._tmp_path.exists(): - self._tmp_path.unlink() - - def _finalize(self) -> None: - """Compare with existing file (if --compare-json) then rename or discard tmp.""" - if self._path in self._files_in_export_dir: - self._files_in_export_dir.remove(self._path) - if self._compare_json: - existing_hash = hashlib.blake2b(self._path.read_bytes()).hexdigest() - new_hash = hashlib.blake2b(self._tmp_path.read_bytes()).hexdigest() - if existing_hash == new_hash: - self._tmp_path.unlink() - return - self._tmp_path.rename(self._path) - - def __enter__(self) -> "StreamingManifestWriter": - self.open() - return self - - def __exit__(self, exc_type, exc_val, exc_tb) -> None: - if exc_type is not None: - self.discard() - else: - self.close() - - class Command(CryptMixin, PaperlessCommand): help = ( "Decrypt and rename all files in our collection into a given target " @@ -314,20 +230,13 @@ class Command(CryptMixin, PaperlessCommand): self.passphrase: str | None = options.get("passphrase") self.batch_size: int = options["batch_size"] - self.files_in_export_dir: set[Path] = set() self.exported_files: set[str] = set() - # If zipping, save the original target for later and - # get a temporary directory for the target instead - temp_dir = None - self.original_target = self.target - if self.zip_export: - settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True) - temp_dir = tempfile.TemporaryDirectory( - dir=settings.SCRATCH_DIR, - prefix="paperless-export", + if self.zip_export and (self.compare_checksums or self.compare_json): + raise CommandError( + "--compare-checksums and --compare-json have no effect when " + "used with --zip", ) - self.target = Path(temp_dir.name).resolve() if not self.target.exists(): raise CommandError("That path doesn't exist") @@ -338,33 +247,28 @@ class Command(CryptMixin, PaperlessCommand): if not os.access(self.target, os.W_OK): raise CommandError("That path doesn't appear to be writable") - try: - # Prevent any ongoing changes in the documents - with FileLock(settings.MEDIA_LOCK): - self.dump() + sink: ExportSink + if self.zip_export: + sink = ZipExportSink( + self.target, + options["zip_name"], + delete=self.delete, + ) + else: + sink = DirectoryExportSink( + self.target, + compare_checksums=self.compare_checksums, + compare_json=self.compare_json, + delete=self.delete, + ) - # We've written everything to the temporary directory in this case, - # now make an archive in the original target, with all files stored - if self.zip_export and temp_dir is not None: - shutil.make_archive( - self.original_target / options["zip_name"], - format="zip", - root_dir=temp_dir.name, - ) + # Prevent any ongoing changes in the documents while exporting + with FileLock(settings.MEDIA_LOCK), sink: + self.dump(sink) - finally: - # Always cleanup the temporary directory, if one was created - if self.zip_export and temp_dir is not None: - temp_dir.cleanup() - - def dump(self) -> None: - # 1. Take a snapshot of what files exist in the current export folder - for x in self.target.glob("**/*"): - if x.is_file(): - self.files_in_export_dir.add(x.resolve()) - - # 2. Create manifest, containing all correspondents, types, tags, storage paths - # note, documents and ui_settings + def dump(self, sink: ExportSink) -> None: + # 1. Create manifest, containing all correspondents, types, tags, storage + # paths, note, documents and ui_settings _excluded_usernames = ["consumer", "AnonymousUser"] manifest_key_to_object_query: dict[str, QuerySet[Any]] = { "correspondents": Correspondent.objects.all(), @@ -427,13 +331,9 @@ class Command(CryptMixin, PaperlessCommand): document_manifest: list[dict] = [] share_link_bundle_manifest: list[dict] = [] - manifest_path = (self.target / "manifest.json").resolve() - with StreamingManifestWriter( - manifest_path, - compare_json=self.compare_json, - files_in_export_dir=self.files_in_export_dir, - ) as writer: + with sink.stream("manifest.json") as handle: + writer = StreamingManifestWriter(handle) with transaction.atomic(): for key, qs in manifest_key_to_object_query.items(): if key == "documents": @@ -469,9 +369,6 @@ class Command(CryptMixin, PaperlessCommand): self._encrypt_record_inline(record) writer.write_batch(batch) - document_map: dict[int, Document] = { - d.pk: d for d in Document.global_objects.order_by("id") - } share_link_bundle_map: dict[int, ShareLinkBundle] = { b.pk: b for b in ShareLinkBundle.objects.order_by("id").prefetch_related( @@ -479,84 +376,72 @@ class Command(CryptMixin, PaperlessCommand): ) } - # 3. Export files from each document - for index, document_dict in enumerate( - self.track( - document_manifest, - description="Exporting documents...", - total=len(document_manifest), - ), + # 2. Export files from each document + # document_manifest and this stream are both ordered by id from the + # same underlying rows, so zip them in lockstep instead of building + # a dict of every Document instance up front (QuerySetStream keeps + # only one batch of documents resident at a time). + documents_stream = QuerySetStream( + Document.global_objects.order_by("id"), + chunk_size=self.batch_size, + ) + for document_dict, document in self.track( + zip(document_manifest, documents_stream, strict=True), + description="Exporting documents...", + total=len(document_manifest), ): - document = document_map[document_dict["pk"]] + # Both document_manifest and documents_stream come from the same + # Document.global_objects.order_by("id") query, taken while + # MEDIA_LOCK is held, so this should be unreachable -- it guards + # against silent data corruption if that invariant ever breaks. + if document.pk != document_dict["pk"]: # pragma: no cover + raise CommandError( + "Document export ordering mismatch: expected " + f"pk={document_dict['pk']}, got pk={document.pk}. " + "Documents may have changed during export.", + ) - # 3.1. generate a unique filename + # generate a unique filename, then the arcnames for its files base_name = self.generate_base_name(document) - - # 3.2. write filenames into manifest - original_target, thumbnail_target, archive_target = ( + original_arc, thumbnail_arc, archive_arc = ( self.generate_document_targets(document, base_name, document_dict) ) - # 3.3. write files to target folder if not self.data_only: self.copy_document_files( document, - original_target, - thumbnail_target, - archive_target, + sink, + original_arc, + thumbnail_arc, + archive_arc, ) if self.split_manifest: - self._write_split_manifest(document_dict, document, base_name) + self._write_split_manifest(sink, document_dict, document, base_name) else: writer.write_record(document_dict) for bundle_dict in share_link_bundle_manifest: bundle = share_link_bundle_map[bundle_dict["pk"]] - - bundle_target = self.generate_share_link_bundle_target( + bundle_arc = self.generate_share_link_bundle_target( bundle, bundle_dict, ) - - if not self.data_only and bundle_target is not None: - self.copy_share_link_bundle_file(bundle, bundle_target) - + if not self.data_only and bundle_arc is not None: + self.copy_share_link_bundle_file(bundle, sink, bundle_arc) writer.write_record(bundle_dict) - # 4.2 write version information to target folder - extra_metadata_path = (self.target / "metadata.json").resolve() + writer.close() + + # 3. Write version (and crypto params) to metadata.json + # Django stores most crypto values in the field itself; we store + # them once here for the whole export metadata: dict[str, str | int | dict[str, str | int]] = { "version": version.__full_version_str__, } - - # 4.2.1 If needed, write the crypto values into the metadata - # Django stores most of these in the field itself, we store them once here if self.passphrase: metadata.update(self.get_crypt_params()) - - self.check_and_write_json( - metadata, - extra_metadata_path, - ) - - if self.delete: - # 5. Remove files which we did not explicitly export in this run - if not self.zip_export: - for f in self.files_in_export_dir: - f.unlink() - - delete_empty_directories( - f.parent, - self.target, - ) - else: - # 5. Remove anything in the original location (before moving the zip) - for item in self.original_target.glob("*"): - if item.is_dir(): - shutil.rmtree(item) - else: - item.unlink() + sink.add_json(metadata, "metadata.json") def generate_base_name(self, document: Document) -> Path: """ @@ -584,73 +469,69 @@ class Command(CryptMixin, PaperlessCommand): document: Document, base_name: Path, document_dict: dict, - ) -> tuple[Path, Path | None, Path | None]: + ) -> tuple[str, str | None, str | None]: """ - Generates the targets for a given document, including the original file, archive file and thumbnail (depending on settings). + Generates the relative POSIX arcnames for a document's original, thumbnail + and archive files (depending on settings), and records them in the manifest. """ original_name = base_name if self.use_folder_prefix: original_name = Path("originals") / original_name - original_target = (self.target / original_name).resolve() - document_dict[EXPORTER_FILE_NAME] = str(original_name) + original_arc = original_name.as_posix() + document_dict[EXPORTER_FILE_NAME] = original_arc if not self.no_thumbnail: thumbnail_name = base_name.parent / (base_name.stem + "-thumbnail.webp") if self.use_folder_prefix: thumbnail_name = Path("thumbnails") / thumbnail_name - thumbnail_target = (self.target / thumbnail_name).resolve() - document_dict[EXPORTER_THUMBNAIL_NAME] = str(thumbnail_name) + thumbnail_arc = thumbnail_name.as_posix() + document_dict[EXPORTER_THUMBNAIL_NAME] = thumbnail_arc else: - thumbnail_target = None + thumbnail_arc = None if not self.no_archive and document.has_archive_version: archive_name = base_name.parent / (base_name.stem + "-archive.pdf") if self.use_folder_prefix: archive_name = Path("archive") / archive_name - archive_target = (self.target / archive_name).resolve() - document_dict[EXPORTER_ARCHIVE_NAME] = str(archive_name) + archive_arc = archive_name.as_posix() + document_dict[EXPORTER_ARCHIVE_NAME] = archive_arc else: - archive_target = None + archive_arc = None - return original_target, thumbnail_target, archive_target + return original_arc, thumbnail_arc, archive_arc def copy_document_files( self, document: Document, - original_target: Path, - thumbnail_target: Path | None, - archive_target: Path | None, + sink: ExportSink, + original_arc: str, + thumbnail_arc: str | None, + archive_arc: str | None, ) -> None: """ - Copies files from the document storage location to the specified target location. - - If the document is encrypted, the files are decrypted before copying them to the target location. + Hands the document's files to the sink (original, thumbnail, archive). """ - self.check_and_copy( - document.source_path, - document.checksum, - original_target, - ) + sink.add_file(document.source_path, original_arc, checksum=document.checksum) - if thumbnail_target: - self.check_and_copy(document.thumbnail_path, None, thumbnail_target) + if thumbnail_arc: + sink.add_file(document.thumbnail_path, thumbnail_arc) - if archive_target: + if archive_arc: if TYPE_CHECKING: assert isinstance(document.archive_path, Path) - self.check_and_copy( + sink.add_file( document.archive_path, - document.archive_checksum, - archive_target, + archive_arc, + checksum=document.archive_checksum, ) def generate_share_link_bundle_target( self, bundle: ShareLinkBundle, bundle_dict: dict, - ) -> Path | None: + ) -> str | None: """ - Generates the export target for a share link bundle file, when present. + Generates the relative POSIX arcname for a share link bundle file, if any. """ if not bundle.file_path: return None @@ -666,25 +547,22 @@ class Command(CryptMixin, PaperlessCommand): bundle_dict["fields"]["file_path"] = portable_bundle_path.as_posix() bundle_dict[EXPORTER_SHARE_LINK_BUNDLE_NAME] = export_bundle_path.as_posix() - return (self.target / export_bundle_path).resolve() + return export_bundle_path.as_posix() def copy_share_link_bundle_file( self, bundle: ShareLinkBundle, - bundle_target: Path, + sink: ExportSink, + bundle_arc: str, ) -> None: """ - Copies a share link bundle ZIP into the export directory. + Hands a share link bundle ZIP to the sink. """ bundle_source_path = bundle.absolute_file_path if bundle_source_path is None: raise FileNotFoundError(f"Share link bundle {bundle.pk} has no file path") - self.check_and_copy( - bundle_source_path, - None, - bundle_target, - ) + sink.add_file(bundle_source_path, bundle_arc) def _encrypt_record_inline(self, record: dict) -> None: """Encrypt sensitive fields in a single record, if passphrase is set.""" @@ -700,6 +578,7 @@ class Command(CryptMixin, PaperlessCommand): def _write_split_manifest( self, + sink: ExportSink, document_dict: dict, document: Document, base_name: Path, @@ -721,81 +600,4 @@ class Command(CryptMixin, PaperlessCommand): manifest_name = base_name.with_name(f"{base_name.stem}-manifest.json") if self.use_folder_prefix: manifest_name = Path("json") / manifest_name - manifest_name = (self.target / manifest_name).resolve() - manifest_name.parent.mkdir(parents=True, exist_ok=True) - self.check_and_write_json(content, manifest_name) - - def check_and_write_json( - self, - content: list[dict] | dict, - target: Path, - ) -> None: - """ - Writes the source content to the target json file. - If --compare-json arg was used, don't write to target file if - the file exists and checksum is identical to content checksum. - This preserves the file timestamps when no changes are made. - """ - - target = target.resolve() - perform_write = True - if target in self.files_in_export_dir: - self.files_in_export_dir.remove(target) - if self.compare_json: - target_checksum = hashlib.blake2b(target.read_bytes()).hexdigest() - src_str = json.dumps( - content, - cls=DjangoJSONEncoder, - indent=2, - ensure_ascii=False, - ) - src_checksum = hashlib.blake2b(src_str.encode("utf-8")).hexdigest() - if src_checksum == target_checksum: - perform_write = False - - if perform_write: - target.write_text( - json.dumps( - content, - cls=DjangoJSONEncoder, - indent=2, - ensure_ascii=False, - ), - encoding="utf-8", - ) - - def check_and_copy( - self, - source: Path, - source_checksum: str | None, - target: Path, - ) -> None: - """ - Copies the source to the target, if target doesn't exist or the target doesn't seem to match - the source attributes - """ - - target = target.resolve() - if target in self.files_in_export_dir: - self.files_in_export_dir.remove(target) - - perform_copy = False - - if target.exists(): - source_stat = source.stat() - target_stat = target.stat() - if self.compare_checksums and source_checksum: - target_checksum = compute_checksum(target) - perform_copy = target_checksum != source_checksum - elif ( - source_stat.st_mtime != target_stat.st_mtime - or source_stat.st_size != target_stat.st_size - ): - perform_copy = True - else: - # Copy if it does not exist - perform_copy = True - - if perform_copy: - target.parent.mkdir(parents=True, exist_ok=True) - copy_file_with_basic_stats(source, target) + sink.add_json(content, manifest_name.as_posix()) diff --git a/src/documents/tests/export/__init__.py b/src/documents/tests/export/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/src/documents/tests/export/test_sinks.py b/src/documents/tests/export/test_sinks.py new file mode 100644 index 000000000..fe421b234 --- /dev/null +++ b/src/documents/tests/export/test_sinks.py @@ -0,0 +1,327 @@ +import io +import json +import os +import zipfile +from pathlib import Path + +import pytest +from pytest_django.fixtures import SettingsWrapper + +from documents.export.sinks import DirectoryExportSink +from documents.export.sinks import ExportSink +from documents.export.sinks import StreamingManifestWriter +from documents.export.sinks import ZipExportSink +from documents.export.sinks import _dumps + + +@pytest.fixture() +def source_file(tmp_path: Path) -> Path: + src: Path = tmp_path / "src" / "doc.pdf" + src.parent.mkdir(parents=True) + src.write_bytes(b"PDF-CONTENT") + return src + + +class TestDumps: + def test_dumps_is_indented_unicode_json(self) -> None: + result: str = _dumps({"a": "é", "b": 1}) + assert '"é"' in result # ensure_ascii=False keeps unicode literal + assert "\n" in result # indent=2 produces newlines + assert json.loads(result) == {"a": "é", "b": 1} + + +class TestStreamingManifestWriter: + def test_writes_json_array_of_records(self) -> None: + handle: io.StringIO = io.StringIO() + writer: StreamingManifestWriter = StreamingManifestWriter(handle) + writer.write_batch([{"pk": 1}, {"pk": 2}]) + writer.write_record({"pk": 3}) + writer.close() + assert json.loads(handle.getvalue()) == [{"pk": 1}, {"pk": 2}, {"pk": 3}] + + def test_empty_manifest_is_valid_empty_array(self) -> None: + handle: io.StringIO = io.StringIO() + writer: StreamingManifestWriter = StreamingManifestWriter(handle) + writer.close() + assert json.loads(handle.getvalue()) == [] + + +class TestDirectoryExportSink: + def test_add_file_copies_to_relative_arcname( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + with DirectoryExportSink( + target, + compare_checksums=False, + compare_json=False, + delete=False, + ) as sink: + sink.add_file(source_file, "originals/doc.pdf") + assert (target / "originals" / "doc.pdf").read_bytes() == b"PDF-CONTENT" + + def test_add_json_writes_file(self, tmp_path: Path) -> None: + target: Path = tmp_path / "out" + target.mkdir() + with DirectoryExportSink( + target, + compare_checksums=False, + compare_json=False, + delete=False, + ) as sink: + sink.add_json({"version": "x"}, "metadata.json") + assert json.loads((target / "metadata.json").read_text()) == {"version": "x"} + + def test_stream_writes_manifest(self, tmp_path: Path) -> None: + target: Path = tmp_path / "out" + target.mkdir() + with DirectoryExportSink( + target, + compare_checksums=False, + compare_json=False, + delete=False, + ) as sink: + with sink.stream("manifest.json") as handle: + writer: StreamingManifestWriter = StreamingManifestWriter(handle) + writer.write_record({"pk": 1}) + writer.close() + assert json.loads((target / "manifest.json").read_text()) == [{"pk": 1}] + + def test_add_file_skips_when_size_and_mtime_match( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + # Pre-existing target with identical size+mtime but DIFFERENT content: + # if add_file skips (no compare-checksums), the old content survives. + target: Path = tmp_path / "out" + target.mkdir() + existing: Path = target / "originals" / "doc.pdf" + existing.parent.mkdir(parents=True) + # Same byte length as the source but different content + matching mtime, + # so a size/mtime comparison treats it as unchanged and skips the copy. + existing.write_bytes(b"X" * len(b"PDF-CONTENT")) + stat = source_file.stat() + os.utime(existing, (stat.st_atime, stat.st_mtime)) + with DirectoryExportSink( + target, + compare_checksums=False, + compare_json=False, + delete=False, + ) as sink: + sink.add_file(source_file, "originals/doc.pdf", checksum="abc") + assert existing.read_bytes() == b"X" * len(b"PDF-CONTENT") # skipped + + def test_add_file_recopies_when_compare_checksums_differ( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + existing: Path = target / "originals" / "doc.pdf" + existing.parent.mkdir(parents=True) + existing.write_bytes(b"X" * len(b"PDF-CONTENT")) + stat = source_file.stat() + os.utime(existing, (stat.st_atime, stat.st_mtime)) + with DirectoryExportSink( + target, + compare_checksums=True, + compare_json=False, + delete=False, + ) as sink: + # wrong checksum forces recopy despite matching size/mtime + sink.add_file(source_file, "originals/doc.pdf", checksum="not-the-real-sum") + assert existing.read_bytes() == b"PDF-CONTENT" # recopied + + def test_delete_prunes_unwritten_snapshot_files( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + stale: Path = target / "stale.pdf" + stale.write_bytes(b"STALE") + with DirectoryExportSink( + target, + compare_checksums=False, + compare_json=False, + delete=True, + ) as sink: + sink.add_file(source_file, "originals/doc.pdf") + assert not stale.exists() + assert (target / "originals" / "doc.pdf").exists() + + def test_no_delete_keeps_unwritten_files( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + stale: Path = target / "stale.pdf" + stale.write_bytes(b"STALE") + with DirectoryExportSink( + target, + compare_checksums=False, + compare_json=False, + delete=False, + ) as sink: + sink.add_file(source_file, "originals/doc.pdf") + assert stale.exists() + + +class TestZipExportSink: + def test_round_trip_files_json_and_stream( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + with ZipExportSink(target, "export", delete=False) as sink: + sink.add_file(source_file, "originals/doc.pdf") + sink.add_json({"version": "x"}, "metadata.json") + with sink.stream("manifest.json") as handle: + writer = StreamingManifestWriter(handle) + writer.write_record({"pk": 1}) + writer.close() + zip_path: Path = target / "export.zip" + assert zip_path.exists() + assert not (target / "export.zip.tmp").exists() + with zipfile.ZipFile(zip_path) as zf: + names = set(zf.namelist()) + assert {"originals/doc.pdf", "metadata.json", "manifest.json"} <= names + assert zf.read("originals/doc.pdf") == b"PDF-CONTENT" + assert json.loads(zf.read("manifest.json")) == [{"pk": 1}] + + def test_nested_arcname_emits_directory_marker( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + with ZipExportSink(target, "export", delete=False) as sink: + sink.add_file(source_file, "originals/doc.pdf") + with zipfile.ZipFile(target / "export.zip") as zf: + assert "originals/" in zf.namelist() + + def test_flat_arcname_has_no_directory_markers( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + with ZipExportSink(target, "export", delete=False) as sink: + sink.add_file(source_file, "doc.pdf") + with zipfile.ZipFile(target / "export.zip") as zf: + assert all(not n.endswith("/") for n in zf.namelist()) + + def test_exception_leaves_no_zip_and_no_tmp( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + with pytest.raises(RuntimeError): + with ZipExportSink(target, "export", delete=False) as sink: + sink.add_file(source_file, "doc.pdf") + raise RuntimeError("boom") + assert not (target / "export.zip").exists() + assert not (target / "export.zip.tmp").exists() + + def test_exception_inside_stream_cleans_up_manifest_tmp( + self, + tmp_path: Path, + source_file: Path, + settings: SettingsWrapper, + ) -> None: + scratch_dir = tmp_path / "scratch" + settings.SCRATCH_DIR = scratch_dir + target: Path = tmp_path / "out" + target.mkdir() + with pytest.raises(RuntimeError): + with ZipExportSink(target, "export", delete=False) as sink: + sink.add_file(source_file, "doc.pdf") + with sink.stream("manifest.json") as handle: + handle.write("[") + raise RuntimeError("boom") + assert list(scratch_dir.glob("export-manifest-*")) == [] + assert not (target / "export.zip").exists() + assert not (target / "export.zip.tmp").exists() + + def test_abort_after_manifest_written_cleans_up_pending_tmp( + self, + tmp_path: Path, + settings: SettingsWrapper, + ) -> None: + scratch_dir = tmp_path / "scratch" + settings.SCRATCH_DIR = scratch_dir + target: Path = tmp_path / "out" + target.mkdir() + with pytest.raises(RuntimeError): + with ZipExportSink(target, "export", delete=False) as sink: + with sink.stream("manifest.json") as handle: + handle.write("[]") + raise RuntimeError("boom") + assert list(scratch_dir.glob("export-manifest-*")) == [] + assert not (target / "export.zip").exists() + + def test_delete_wipes_destination_on_success( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + (target / "preexisting.txt").write_text("old") + (target / "olddir").mkdir() + with ZipExportSink(target, "export", delete=True) as sink: + sink.add_file(source_file, "doc.pdf") + assert (target / "export.zip").exists() + assert not (target / "preexisting.txt").exists() + assert not (target / "olddir").exists() + + def test_abort_with_delete_does_not_wipe_destination( + self, + tmp_path: Path, + source_file: Path, + ) -> None: + target: Path = tmp_path / "out" + target.mkdir() + (target / "preexisting.txt").write_text("old") + with pytest.raises(RuntimeError): + with ZipExportSink(target, "export", delete=True) as sink: + sink.add_file(source_file, "doc.pdf") + raise RuntimeError("boom") + assert (target / "preexisting.txt").exists() + assert not (target / "export.zip").exists() + + +class TestStreamContract: + @pytest.fixture(params=["dir", "zip"]) + def sink(self, request: pytest.FixtureRequest, tmp_path: Path) -> ExportSink: + target: Path = tmp_path / "out" + target.mkdir() + if request.param == "dir": + return DirectoryExportSink( + target, + compare_checksums=False, + compare_json=False, + delete=False, + ) + return ZipExportSink(target, "export", delete=False) + + def test_second_concurrent_stream_is_rejected(self, sink: ExportSink) -> None: + with sink: + with sink.stream("manifest.json"): + with pytest.raises(RuntimeError, match="already open"): + with sink.stream("other.json"): + pass diff --git a/src/documents/tests/test_management_exporter.py b/src/documents/tests/test_management_exporter.py index 4ee7677ca..44456535a 100644 --- a/src/documents/tests/test_management_exporter.py +++ b/src/documents/tests/test_management_exporter.py @@ -426,7 +426,7 @@ class TestExportImport( st_mtime_1 = (self.target / "manifest.json").stat().st_mtime with mock.patch( - "documents.management.commands.document_exporter.copy_file_with_basic_stats", + "documents.export.sinks.copy_file_with_basic_stats", ) as m: self._do_export() m.assert_not_called() @@ -437,7 +437,7 @@ class TestExportImport( Path(self.d1.source_path).touch() with mock.patch( - "documents.management.commands.document_exporter.copy_file_with_basic_stats", + "documents.export.sinks.copy_file_with_basic_stats", ) as m: self._do_export() self.assertEqual(m.call_count, 1) @@ -464,7 +464,7 @@ class TestExportImport( self.assertIsFile(self.target / "manifest.json") with mock.patch( - "documents.management.commands.document_exporter.copy_file_with_basic_stats", + "documents.export.sinks.copy_file_with_basic_stats", ) as m: self._do_export() m.assert_not_called() @@ -475,7 +475,7 @@ class TestExportImport( self.d2.save() with mock.patch( - "documents.management.commands.document_exporter.copy_file_with_basic_stats", + "documents.export.sinks.copy_file_with_basic_stats", ) as m: self._do_export(compare_checksums=True) self.assertEqual(m.call_count, 1) @@ -1058,6 +1058,26 @@ class TestExportImport( self.assertEqual(Document.objects.all().count(), 4) + def test_zip_with_compare_flags_raises(self) -> None: + """ + GIVEN: + - A request to export to a zip file + WHEN: + - --compare-checksums or --compare-json is also passed + THEN: + - A CommandError is raised (the flags are no-ops in zip mode) + """ + for flag in ("--compare-checksums", "--compare-json"): + with self.subTest(flag=flag): + with self.assertRaises(CommandError): + call_command( + "document_exporter", + self.target, + "--zip", + flag, + skip_checks=True, + ) + @pytest.mark.management class TestCryptExportImport(