Coverage for documents/export/sinks.py: 0%
210 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
1from __future__ import annotations
3import abc
4import hashlib
5import json
6import os
7import shutil
8import tempfile
9import zipfile
10from contextlib import AbstractContextManager
11from contextlib import contextmanager
12from pathlib import Path
13from pathlib import PurePosixPath
14from typing import TYPE_CHECKING
16from django.conf import settings
17from django.core.serializers.json import DjangoJSONEncoder
19from documents.file_handling import delete_empty_directories
20from documents.utils import compute_checksum
21from documents.utils import copy_file_with_basic_stats
23if TYPE_CHECKING:
24 from collections.abc import Iterator
25 from typing import TextIO
28def _dumps(content: list | dict) -> str:
29 """Serialize export JSON consistently across all sinks."""
30 return json.dumps(content, cls=DjangoJSONEncoder, indent=2, ensure_ascii=False)
33class StreamingManifestWriter:
34 """Incrementally writes a JSON array to a text handle, one record at a time.
36 Knows nothing about folders or zips: it writes the array framing and records
37 to whatever handle the sink's ``stream()`` yields. The sink owns the handle's
38 lifecycle (atomic rename, compare, spooling).
39 """
41 def __init__(self, handle: TextIO) -> None:
42 self._file = handle
43 self._first = True
44 self._file.write("[")
46 def write_record(self, record: dict) -> None:
47 if not self._first:
48 self._file.write(",\n")
49 else:
50 self._first = False
51 self._file.write(_dumps(record))
53 def write_batch(self, records: list[dict]) -> None:
54 for record in records:
55 self.write_record(record)
57 def close(self) -> None:
58 """Write the closing bracket. Does NOT close the handle (the sink owns it)."""
59 self._file.write("\n]")
62class ExportSink(AbstractContextManager, abc.ABC):
63 """Destination for a document export.
65 The command declares export contents via three verbs; the sink decides how to
66 persist each. ``arcname`` is always a relative POSIX path
67 (e.g. ``"manifest.json"``, ``"originals/foo.pdf"``).
69 Contract:
70 * At most one ``stream()`` open at a time (it is the manifest);
71 ``add_file``/``add_json`` may be called while it is open.
72 * Context-manager: normal exit finalizes, an exception aborts. No partial or
73 failed run leaves a complete-looking artifact.
74 """
76 @abc.abstractmethod
77 def add_file(
78 self,
79 source: Path,
80 arcname: str,
81 *,
82 checksum: str | None = None,
83 ) -> None: ...
85 @abc.abstractmethod
86 def add_json(self, content: list | dict, arcname: str) -> None: ...
88 @abc.abstractmethod
89 def stream(self, arcname: str) -> AbstractContextManager[TextIO]: ...
91 def _open(self) -> None:
92 """Hook called on context entry. Override as needed."""
94 @abc.abstractmethod
95 def _finalize(self) -> None:
96 """Commit on clean exit."""
98 @abc.abstractmethod
99 def _abort(self) -> None:
100 """Roll back on exception."""
102 def __enter__(self) -> ExportSink:
103 self._open()
104 return self
106 def __exit__(self, exc_type, exc_val, exc_tb) -> None:
107 if exc_type is not None:
108 self._abort()
109 else:
110 self._finalize()
113class DirectoryExportSink(ExportSink):
114 """Writes loose files into a target directory, with incremental sync.
116 Owns the snapshot/skip/compare/prune machinery that used to live in the
117 command (``files_in_export_dir``, ``check_and_copy``, ``check_and_write_json``,
118 and the ``--delete`` pass).
119 """
121 def __init__(
122 self,
123 target: Path,
124 *,
125 compare_checksums: bool,
126 compare_json: bool,
127 delete: bool,
128 ) -> None:
129 self._target = target.resolve()
130 self._compare_checksums = compare_checksums
131 self._compare_json = compare_json
132 self._delete = delete
133 self._snapshot: set[Path] = set()
134 self._stream_open = False
136 def _open(self) -> None:
137 for x in self._target.glob("**/*"):
138 if x.is_file():
139 self._snapshot.add(x.resolve())
141 def add_file(
142 self,
143 source: Path,
144 arcname: str,
145 *,
146 checksum: str | None = None,
147 ) -> None:
148 target = (self._target / arcname).resolve()
149 self._snapshot.discard(target)
150 perform_copy = False
151 if target.exists():
152 source_stat = source.stat()
153 target_stat = target.stat()
154 if self._compare_checksums and checksum:
155 perform_copy = compute_checksum(target) != checksum
156 elif (
157 source_stat.st_mtime != target_stat.st_mtime
158 or source_stat.st_size != target_stat.st_size
159 ):
160 perform_copy = True
161 else:
162 perform_copy = True
163 if perform_copy:
164 target.parent.mkdir(parents=True, exist_ok=True)
165 copy_file_with_basic_stats(source, target)
167 @staticmethod
168 def _content_unchanged(target: Path, new_bytes: bytes) -> bool:
169 """True if ``target`` already holds byte-identical content (BLAKE2b)."""
170 return (
171 hashlib.blake2b(target.read_bytes()).hexdigest()
172 == hashlib.blake2b(new_bytes).hexdigest()
173 )
175 def add_json(self, content: list | dict, arcname: str) -> None:
176 target = (self._target / arcname).resolve()
177 json_str = _dumps(content)
178 perform_write = True
179 if target in self._snapshot:
180 self._snapshot.discard(target)
181 if self._compare_json and self._content_unchanged(
182 target,
183 json_str.encode("utf-8"),
184 ):
185 perform_write = False
186 if perform_write:
187 target.parent.mkdir(parents=True, exist_ok=True)
188 target.write_text(json_str, encoding="utf-8")
190 @contextmanager
191 def stream(self, arcname: str) -> Iterator[TextIO]:
192 if self._stream_open:
193 raise RuntimeError("A stream is already open on this sink")
194 target = (self._target / arcname).resolve()
195 tmp = target.with_suffix(target.suffix + ".tmp")
196 target.parent.mkdir(parents=True, exist_ok=True)
197 handle = tmp.open("w", encoding="utf-8")
198 self._stream_open = True
199 try:
200 yield handle
201 except BaseException:
202 handle.close()
203 tmp.unlink(missing_ok=True)
204 raise
205 else:
206 handle.close()
207 self._commit_streamed_file(target, tmp)
208 finally:
209 self._stream_open = False
211 def _commit_streamed_file(self, target: Path, tmp: Path) -> None:
212 if target in self._snapshot:
213 self._snapshot.discard(target)
214 if self._compare_json and self._content_unchanged(
215 target,
216 tmp.read_bytes(),
217 ):
218 tmp.unlink()
219 return
220 tmp.rename(target)
222 def _finalize(self) -> None:
223 if self._delete:
224 for f in self._snapshot:
225 if not f.is_relative_to(self._target): # pragma: no cover
226 # Defense in depth: a symlink inside the export dir can
227 # resolve outside of it; never delete outside the target.
228 continue
229 f.unlink()
230 delete_empty_directories(f.parent, self._target)
232 def _abort(self) -> None:
233 # Folder mode is in-place/incremental: streamed .tmp files are already
234 # cleaned in stream(); leave everything else intact and skip the prune.
235 return None
238class ZipExportSink(ExportSink):
239 """Writes a single zip archive, produced atomically only on success.
241 Builds into ``<target>/<zip_name>.zip.tmp`` and renames to ``.zip`` on clean
242 finalize. The manifest stream is spooled to a temp file in SCRATCH_DIR and
243 added as an entry at finalize (a zip entry cannot be interleaved with others).
244 """
246 def __init__(
247 self,
248 target: Path,
249 zip_name: str,
250 *,
251 delete: bool = False,
252 compression: int = zipfile.ZIP_DEFLATED,
253 compresslevel: int | None = None,
254 ) -> None:
255 self._target = target.resolve()
256 self._zip_path = (self._target / zip_name).with_suffix(".zip")
257 self._tmp_path = self._zip_path.with_name(self._zip_path.name + ".tmp")
258 self._delete = delete
259 self._compression = compression
260 self._compresslevel = compresslevel
261 self._zip: zipfile.ZipFile | None = None
262 self._dirs: set[str] = set()
263 self._pending_manifest: tuple[Path, str] | None = None
264 self._stream_open = False
266 def _open(self) -> None:
267 settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True)
268 self._zip = zipfile.ZipFile(
269 self._tmp_path,
270 "w",
271 compression=self._compression,
272 compresslevel=self._compresslevel,
273 allowZip64=True,
274 )
276 def _ensure_dirs(self, arcname: str) -> None:
277 assert self._zip is not None
278 dir_arc = ""
279 for part in PurePosixPath(arcname).parts[:-1]:
280 dir_arc += f"{part}/"
281 if dir_arc not in self._dirs:
282 self._dirs.add(dir_arc)
283 self._zip.mkdir(dir_arc)
285 def add_file(
286 self,
287 source: Path,
288 arcname: str,
289 *,
290 checksum: str | None = None,
291 ) -> None:
292 assert self._zip is not None
293 self._ensure_dirs(arcname)
294 self._zip.write(source, arcname=arcname)
296 def add_json(self, content: list | dict, arcname: str) -> None:
297 assert self._zip is not None
298 self._ensure_dirs(arcname)
299 self._zip.writestr(arcname, _dumps(content))
301 @contextmanager
302 def stream(self, arcname: str) -> Iterator[TextIO]:
303 if self._stream_open:
304 raise RuntimeError("A stream is already open on this sink")
305 settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True)
306 fd, tmp_name = tempfile.mkstemp(
307 dir=settings.SCRATCH_DIR,
308 prefix="export-manifest-",
309 suffix=".json",
310 )
311 tmp = Path(tmp_name)
312 handle = os.fdopen(fd, "w", encoding="utf-8")
313 self._stream_open = True
314 try:
315 yield handle
316 except BaseException:
317 handle.close()
318 tmp.unlink(missing_ok=True)
319 raise
320 else:
321 handle.close()
322 self._pending_manifest = (tmp, arcname)
323 finally:
324 self._stream_open = False
326 def _finalize(self) -> None:
327 assert self._zip is not None
328 if self._pending_manifest is not None:
329 tmp, arcname = self._pending_manifest
330 self._ensure_dirs(arcname)
331 self._zip.write(tmp, arcname=arcname)
332 tmp.unlink(missing_ok=True)
333 self._pending_manifest = None
334 self._zip.close()
335 self._zip = None
336 if self._delete:
337 self._wipe_destination()
338 self._tmp_path.replace(self._zip_path)
340 def _wipe_destination(self) -> None:
341 skip = {self._zip_path.resolve(), self._tmp_path.resolve()}
342 for item in self._target.glob("*"):
343 if item.resolve() in skip:
344 continue
345 if item.is_dir():
346 shutil.rmtree(item)
347 else:
348 item.unlink()
350 def _abort(self) -> None:
351 if self._zip is not None:
352 self._zip.close()
353 self._zip = None
354 self._tmp_path.unlink(missing_ok=True)
355 if self._pending_manifest is not None:
356 self._pending_manifest[0].unlink(missing_ok=True)
357 self._pending_manifest = None