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

1from __future__ import annotations 

2 

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 

15 

16from django.conf import settings 

17from django.core.serializers.json import DjangoJSONEncoder 

18 

19from documents.file_handling import delete_empty_directories 

20from documents.utils import compute_checksum 

21from documents.utils import copy_file_with_basic_stats 

22 

23if TYPE_CHECKING: 

24 from collections.abc import Iterator 

25 from typing import TextIO 

26 

27 

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) 

31 

32 

33class StreamingManifestWriter: 

34 """Incrementally writes a JSON array to a text handle, one record at a time. 

35 

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 """ 

40 

41 def __init__(self, handle: TextIO) -> None: 

42 self._file = handle 

43 self._first = True 

44 self._file.write("[") 

45 

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)) 

52 

53 def write_batch(self, records: list[dict]) -> None: 

54 for record in records: 

55 self.write_record(record) 

56 

57 def close(self) -> None: 

58 """Write the closing bracket. Does NOT close the handle (the sink owns it).""" 

59 self._file.write("\n]") 

60 

61 

62class ExportSink(AbstractContextManager, abc.ABC): 

63 """Destination for a document export. 

64 

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"``). 

68 

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 """ 

75 

76 @abc.abstractmethod 

77 def add_file( 

78 self, 

79 source: Path, 

80 arcname: str, 

81 *, 

82 checksum: str | None = None, 

83 ) -> None: ... 

84 

85 @abc.abstractmethod 

86 def add_json(self, content: list | dict, arcname: str) -> None: ... 

87 

88 @abc.abstractmethod 

89 def stream(self, arcname: str) -> AbstractContextManager[TextIO]: ... 

90 

91 def _open(self) -> None: 

92 """Hook called on context entry. Override as needed.""" 

93 

94 @abc.abstractmethod 

95 def _finalize(self) -> None: 

96 """Commit on clean exit.""" 

97 

98 @abc.abstractmethod 

99 def _abort(self) -> None: 

100 """Roll back on exception.""" 

101 

102 def __enter__(self) -> ExportSink: 

103 self._open() 

104 return self 

105 

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() 

111 

112 

113class DirectoryExportSink(ExportSink): 

114 """Writes loose files into a target directory, with incremental sync. 

115 

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 """ 

120 

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 

135 

136 def _open(self) -> None: 

137 for x in self._target.glob("**/*"): 

138 if x.is_file(): 

139 self._snapshot.add(x.resolve()) 

140 

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) 

166 

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 ) 

174 

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") 

189 

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 

210 

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) 

221 

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) 

231 

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 

236 

237 

238class ZipExportSink(ExportSink): 

239 """Writes a single zip archive, produced atomically only on success. 

240 

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 """ 

245 

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 

265 

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 ) 

275 

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) 

284 

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) 

295 

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)) 

300 

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 

325 

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) 

339 

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() 

349 

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