Coverage for documents/consumer.py: 21%

449 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 09:07 +0000

1import datetime 

2import logging 

3import os 

4import shutil 

5import tempfile 

6from enum import StrEnum 

7from pathlib import Path 

8from typing import TYPE_CHECKING 

9from typing import Final 

10 

11import magic 

12from django.conf import settings 

13from django.contrib.auth.models import User 

14from django.db import transaction 

15from django.db.models import Max 

16from django.db.models import Q 

17from django.utils import timezone 

18from filelock import FileLock 

19from rest_framework.reverse import reverse 

20 

21from documents.barcodes import read_barcode_values 

22from documents.classifier import load_classifier 

23from documents.data_models import ConsumableDocument 

24from documents.data_models import ConsumeFileSuccessResult 

25from documents.data_models import DocumentMetadataOverrides 

26from documents.file_handling import create_source_path_directory 

27from documents.file_handling import generate_filename 

28from documents.file_handling import generate_unique_filename 

29from documents.file_handling import validate_path_in_root 

30from documents.loggers import LoggingMixin 

31from documents.models import Correspondent 

32from documents.models import CustomField 

33from documents.models import CustomFieldInstance 

34from documents.models import Document 

35from documents.models import DocumentBarcode 

36from documents.models import DocumentType 

37from documents.models import StoragePath 

38from documents.models import Tag 

39from documents.models import WorkflowTrigger 

40from documents.parsers import ParseError 

41from documents.permissions import set_permissions_for_object 

42from documents.plugins.base import AlwaysRunPluginMixin 

43from documents.plugins.base import ConsumeTaskPlugin 

44from documents.plugins.base import NoCleanupPluginMixin 

45from documents.plugins.base import NoSetupPluginMixin 

46from documents.plugins.date_parsing import get_date_parser 

47from documents.plugins.helpers import ProgressManager 

48from documents.plugins.helpers import ProgressStatusOptions 

49from documents.signals import document_consumption_finished 

50from documents.signals import document_consumption_started 

51from documents.signals import document_updated 

52from documents.signals.handlers import run_workflows 

53from documents.templating.workflows import parse_w_workflow_placeholders 

54from documents.utils import compute_checksum 

55from documents.utils import copy_basic_file_stats 

56from documents.utils import copy_file_with_basic_stats 

57from documents.utils import run_subprocess 

58from paperless.config import BarcodeConfig 

59from paperless.config import OcrConfig 

60from paperless.config import RemoteOCRConfig 

61from paperless.models import ArchiveFileGenerationChoices 

62from paperless.parsers import ParserContext 

63from paperless.parsers import ParserProtocol 

64from paperless.parsers.registry import get_parser_registry 

65from paperless.parsers.utils import pdf_born_digital_text 

66 

67LOGGING_NAME: Final[str] = "paperless.consumer" 

68 

69 

70class WorkflowTriggerPlugin( 

71 NoCleanupPluginMixin, 

72 NoSetupPluginMixin, 

73 AlwaysRunPluginMixin, 

74 ConsumeTaskPlugin, 

75): 

76 NAME: str = "WorkflowTriggerPlugin" 

77 

78 def run(self) -> str | None: 

79 """ 

80 Get overrides from matching workflows 

81 """ 

82 overrides, msg = run_workflows( 

83 trigger_type=WorkflowTrigger.WorkflowTriggerType.CONSUMPTION, 

84 document=self.input_doc, 

85 logging_group=None, 

86 overrides=DocumentMetadataOverrides(), 

87 ) 

88 if overrides: 

89 self.metadata.update(overrides) 

90 return msg 

91 

92 

93class ConsumerError(Exception): 

94 pass 

95 

96 

97class ConsumeFileDuplicateError(ConsumerError): 

98 """Raised when a file is rejected because it duplicates an existing document.""" 

99 

100 def __init__(self, message: str, duplicate_id: int, *, in_trash: bool) -> None: 

101 super().__init__(message) 

102 self.duplicate_id = duplicate_id 

103 self.in_trash = in_trash 

104 

105 

106class ConsumerStatusShortMessage(StrEnum): 

107 DOCUMENT_ALREADY_EXISTS = "document_already_exists" 

108 DOCUMENT_ALREADY_EXISTS_IN_TRASH = "document_already_exists_in_trash" 

109 ASN_ALREADY_EXISTS = "asn_already_exists" 

110 ASN_ALREADY_EXISTS_IN_TRASH = "asn_already_exists_in_trash" 

111 ASN_RANGE = "asn_value_out_of_range" 

112 FILE_NOT_FOUND = "file_not_found" 

113 PRE_CONSUME_SCRIPT_NOT_FOUND = "pre_consume_script_not_found" 

114 PRE_CONSUME_SCRIPT_ERROR = "pre_consume_script_error" 

115 POST_CONSUME_SCRIPT_NOT_FOUND = "post_consume_script_not_found" 

116 POST_CONSUME_SCRIPT_ERROR = "post_consume_script_error" 

117 NEW_FILE = "new_file" 

118 UNSUPPORTED_TYPE = "unsupported_type" 

119 PARSING_DOCUMENT = "parsing_document" 

120 GENERATING_THUMBNAIL = "generating_thumbnail" 

121 PARSE_DATE = "parse_date" 

122 SAVE_DOCUMENT = "save_document" 

123 FINISHED = "finished" 

124 FAILED = "failed" 

125 

126 

127def should_produce_archive( 

128 parser: "ParserProtocol", 

129 mime_type: str, 

130 document_path: Path, 

131 log: logging.Logger | None = None, 

132) -> bool: 

133 """Return True if a PDF/A archive should be produced for this document. 

134 

135 IMPORTANT: *parser* must be an instantiated parser, not the class. 

136 ``requires_pdf_rendition`` and ``can_produce_archive`` are instance 

137 ``@property`` methods — accessing them on the class returns the descriptor 

138 (always truthy). 

139 """ 

140 _log = log or logging.getLogger(LOGGING_NAME) 

141 

142 # Must produce a PDF so the frontend can display the original format at all. 

143 if parser.requires_pdf_rendition: 

144 _log.debug("Archive: yes - parser requires PDF rendition for frontend display") 

145 return True 

146 

147 # Parser cannot produce an archive (e.g. TextDocumentParser). 

148 if not parser.can_produce_archive: 

149 _log.debug("Archive: no - parser cannot produce archives") 

150 return False 

151 

152 generation = OcrConfig().archive_file_generation 

153 

154 if generation == ArchiveFileGenerationChoices.ALWAYS: 

155 _log.debug("Archive: yes - ARCHIVE_FILE_GENERATION=always") 

156 return True 

157 if generation == ArchiveFileGenerationChoices.NEVER: 

158 _log.debug("Archive: no - ARCHIVE_FILE_GENERATION=never") 

159 return False 

160 

161 # auto: produce archives for scanned/image documents; skip for born-digital PDFs. 

162 if mime_type.startswith("image/"): 

163 _log.debug("Archive: yes - image document, ARCHIVE_FILE_GENERATION=auto") 

164 return True 

165 if mime_type == "application/pdf": 

166 text, born_digital = pdf_born_digital_text(document_path, log=_log) 

167 text_length = len(text) if text else 0 

168 if born_digital: 

169 _log.debug( 

170 "Archive: no - born-digital PDF (text_length=%d)," 

171 " ARCHIVE_FILE_GENERATION=auto", 

172 text_length, 

173 ) 

174 return False 

175 _log.debug( 

176 "Archive: yes - scanned/textless PDF (text_length=%d)," 

177 " ARCHIVE_FILE_GENERATION=auto", 

178 text_length, 

179 ) 

180 return True 

181 _log.debug( 

182 "Archive: no - MIME type %r not eligible for auto archive generation", 

183 mime_type, 

184 ) 

185 return False 

186 

187 

188class ConsumerPluginMixin: 

189 if TYPE_CHECKING: 189 ↛ 190line 189 didn't jump to line 190 because the condition on line 189 was never true

190 from logging import Logger 

191 from logging import LoggerAdapter 

192 

193 log: "LoggerAdapter" # type: ignore[type-arg] 

194 

195 def __init__( 

196 self, 

197 input_doc: ConsumableDocument, 

198 metadata: DocumentMetadataOverrides, 

199 status_mgr: ProgressManager, 

200 base_tmp_dir: Path, 

201 task_id: str, 

202 ) -> None: 

203 super().__init__(input_doc, metadata, status_mgr, base_tmp_dir, task_id) 

204 

205 self.renew_logging_group() 

206 

207 self.filename = self.metadata.filename or self.input_doc.original_file.name 

208 

209 def _send_progress( 

210 self, 

211 current_progress: int, 

212 max_progress: int, 

213 status: ProgressStatusOptions, 

214 message: ConsumerStatusShortMessage | str | None = None, 

215 document_id=None, 

216 ) -> None: # pragma: no cover 

217 self.status_mgr.send_progress( 

218 status, 

219 message, 

220 current_progress, 

221 max_progress, 

222 document_id=document_id, 

223 owner_id=self.metadata.owner_id if self.metadata.owner_id else None, 

224 users_can_view=(self.metadata.view_users or []) 

225 + (self.metadata.change_users or []), 

226 groups_can_view=(self.metadata.view_groups or []) 

227 + (self.metadata.change_groups or []), 

228 ) 

229 

230 def _fail( 

231 self, 

232 message: ConsumerStatusShortMessage | str, 

233 log_message: str | None = None, 

234 exc_info=None, 

235 exception: Exception | None = None, 

236 ): 

237 self._send_progress(100, 100, ProgressStatusOptions.FAILED, message) 

238 self.log.error(log_message or message, exc_info=exc_info) 

239 raise ConsumerError(f"{self.filename}: {log_message or message}") from exception 

240 

241 

242class ConsumerPlugin( 

243 AlwaysRunPluginMixin, 

244 NoSetupPluginMixin, 

245 NoCleanupPluginMixin, 

246 LoggingMixin, 

247 ConsumerPluginMixin, 

248 ConsumeTaskPlugin, 

249): 

250 logging_name = LOGGING_NAME 

251 

252 def _create_version_from_root( 

253 self, 

254 root_doc: Document, 

255 *, 

256 text: str | None, 

257 page_count: int | None, 

258 mime_type: str, 

259 ) -> Document: 

260 self.log.debug("Saving record for updated version to database") 

261 root_doc_frozen = Document.objects.select_for_update().get(pk=root_doc.pk) 

262 next_version_index = ( 

263 Document.global_objects.filter( 

264 root_document_id=root_doc_frozen.pk, 

265 ).aggregate( 

266 max_index=Max("version_index"), 

267 )["max_index"] 

268 or 0 

269 ) 

270 file_for_checksum = ( 

271 self.unmodified_original 

272 if self.unmodified_original is not None 

273 else self.working_copy 

274 ) 

275 version_doc = Document( 

276 root_document=root_doc_frozen, 

277 version_index=next_version_index + 1, 

278 checksum=compute_checksum(file_for_checksum), 

279 content=text or "", 

280 page_count=page_count, 

281 mime_type=mime_type, 

282 original_filename=self.filename, 

283 owner_id=root_doc_frozen.owner_id, 

284 created=root_doc_frozen.created, 

285 title=root_doc_frozen.title, 

286 added=timezone.now(), 

287 modified=timezone.now(), 

288 ) 

289 if self.metadata.version_label is not None: 

290 version_doc.version_label = self.metadata.version_label 

291 return version_doc 

292 

293 def run_pre_consume_script(self) -> None: 

294 """ 

295 If one is configured and exists, run the pre-consume script and 

296 handle its output and/or errors 

297 """ 

298 if not settings.PRE_CONSUME_SCRIPT: 

299 return 

300 

301 if not Path(settings.PRE_CONSUME_SCRIPT).is_file(): 

302 self._fail( 

303 ConsumerStatusShortMessage.PRE_CONSUME_SCRIPT_NOT_FOUND, 

304 f"Configured pre-consume script " 

305 f"{settings.PRE_CONSUME_SCRIPT} does not exist.", 

306 ) 

307 

308 self.log.info(f"Executing pre-consume script {settings.PRE_CONSUME_SCRIPT}") 

309 

310 working_file_path = str(self.working_copy) 

311 original_file_path = str(self.input_doc.original_file) 

312 

313 script_env = os.environ.copy() 

314 script_env["DOCUMENT_SOURCE_PATH"] = original_file_path 

315 script_env["DOCUMENT_WORKING_PATH"] = working_file_path 

316 script_env["TASK_ID"] = self.task_id or "" 

317 

318 try: 

319 run_subprocess( 

320 [ 

321 settings.PRE_CONSUME_SCRIPT, 

322 ], 

323 script_env, 

324 self.log, 

325 ) 

326 

327 except Exception as e: 

328 self._fail( 

329 ConsumerStatusShortMessage.PRE_CONSUME_SCRIPT_ERROR, 

330 f"Error while executing pre-consume script: {e}", 

331 exc_info=True, 

332 exception=e, 

333 ) 

334 

335 def run_post_consume_script(self, document: Document) -> None: 

336 """ 

337 If one is configured and exists, run the pre-consume script and 

338 handle its output and/or errors 

339 """ 

340 if not settings.POST_CONSUME_SCRIPT: 

341 return 

342 

343 if not Path(settings.POST_CONSUME_SCRIPT).is_file(): 

344 self._fail( 

345 ConsumerStatusShortMessage.POST_CONSUME_SCRIPT_NOT_FOUND, 

346 f"Configured post-consume script " 

347 f"{settings.POST_CONSUME_SCRIPT} does not exist.", 

348 ) 

349 

350 self.log.info( 

351 f"Executing post-consume script {settings.POST_CONSUME_SCRIPT}", 

352 ) 

353 

354 script_env = os.environ.copy() 

355 

356 script_env["DOCUMENT_ID"] = str(document.pk) 

357 script_env["DOCUMENT_TYPE"] = str(document.document_type) 

358 script_env["DOCUMENT_CREATED"] = str(document.created) 

359 script_env["DOCUMENT_MODIFIED"] = str(document.modified) 

360 script_env["DOCUMENT_ADDED"] = str(document.added) 

361 script_env["DOCUMENT_FILE_NAME"] = document.get_public_filename() 

362 script_env["DOCUMENT_SOURCE_PATH"] = os.path.normpath(document.source_path) 

363 script_env["DOCUMENT_ARCHIVE_PATH"] = os.path.normpath( 

364 str(document.archive_path), 

365 ) 

366 script_env["DOCUMENT_THUMBNAIL_PATH"] = os.path.normpath( 

367 document.thumbnail_path, 

368 ) 

369 script_env["DOCUMENT_DOWNLOAD_URL"] = reverse( 

370 "document-download", 

371 kwargs={"pk": document.pk}, 

372 ) 

373 script_env["DOCUMENT_THUMBNAIL_URL"] = reverse( 

374 "document-thumb", 

375 kwargs={"pk": document.pk}, 

376 ) 

377 script_env["DOCUMENT_OWNER"] = ( 

378 document.owner.get_username() if document.owner else "" 

379 ) 

380 script_env["DOCUMENT_CORRESPONDENT"] = str(document.correspondent) 

381 script_env["DOCUMENT_TAGS"] = str( 

382 ",".join(document.tags.all().values_list("name", flat=True)), 

383 ) 

384 script_env["DOCUMENT_ORIGINAL_FILENAME"] = str(document.original_filename) 

385 script_env["TASK_ID"] = self.task_id or "" 

386 

387 try: 

388 run_subprocess( 

389 [ 

390 settings.POST_CONSUME_SCRIPT, 

391 ], 

392 script_env, 

393 self.log, 

394 ) 

395 

396 except Exception as e: 

397 self._fail( 

398 ConsumerStatusShortMessage.POST_CONSUME_SCRIPT_ERROR, 

399 f"Error while executing post-consume script: {e}", 

400 exc_info=True, 

401 exception=e, 

402 ) 

403 

404 def run(self) -> "ConsumeFileSuccessResult": 

405 """ 

406 Return the document object if it was successfully created. 

407 """ 

408 

409 # Preflight has already run including progress update to 0% 

410 self.log.info(f"Consuming {self.filename}") 

411 

412 # For the actual work, copy the file into a tempdir 

413 with tempfile.TemporaryDirectory( 

414 prefix="paperless-ngx", 

415 dir=settings.SCRATCH_DIR, 

416 ) as tmpdir: 

417 self.working_copy = Path(tmpdir) / Path(self.filename) 

418 copy_file_with_basic_stats(self.input_doc.original_file, self.working_copy) 

419 self.unmodified_original = None 

420 

421 # Determine the parser class. 

422 

423 mime_type = magic.from_file(self.working_copy, mime=True) 

424 

425 self.log.debug(f"Detected mime type: {mime_type}") 

426 

427 if ( 

428 Path(self.filename).suffix.lower() == ".pdf" 

429 and mime_type in settings.CONSUMER_PDF_RECOVERABLE_MIME_TYPES 

430 ): 

431 try: 

432 # The file might be a pdf, but the mime type is wrong. 

433 # Try to clean with qpdf 

434 self.log.debug( 

435 "Detected possible PDF with wrong mime type, trying to clean with qpdf", 

436 ) 

437 run_subprocess( 

438 [ 

439 "qpdf", 

440 "--replace-input", 

441 self.working_copy, 

442 ], 

443 logger=self.log, 

444 ) 

445 mime_type = magic.from_file(self.working_copy, mime=True) 

446 self.log.debug(f"Detected mime type after qpdf: {mime_type}") 

447 # Save the original file for later 

448 self.unmodified_original = ( 

449 Path(tmpdir) / Path("uo") / Path(self.filename) 

450 ) 

451 self.unmodified_original.parent.mkdir(exist_ok=True) 

452 copy_file_with_basic_stats( 

453 self.input_doc.original_file, 

454 self.unmodified_original, 

455 ) 

456 except Exception as e: 

457 self.log.error(f"Error attempting to clean PDF: {e}") 

458 

459 # Workflows have already run at this point, so the metadata knows 

460 # whether this document was singled out for remote OCR 

461 allow_remote = ( 

462 self.metadata.remote_ocr or RemoteOCRConfig().remote_ocr_by_default 

463 ) 

464 

465 # Based on the mime type, get the parser for that type 

466 parser_class: type[ParserProtocol] | None = ( 

467 get_parser_registry().get_parser_for_file( 

468 mime_type, 

469 self.filename, 

470 self.working_copy, 

471 allow_remote=allow_remote, 

472 ) 

473 ) 

474 if not parser_class: 

475 self._fail( 

476 ConsumerStatusShortMessage.UNSUPPORTED_TYPE, 

477 f"Unsupported mime type {mime_type}", 

478 ) 

479 

480 if self.metadata.remote_ocr and not getattr( 

481 parser_class, 

482 "uses_remote_service", 

483 False, 

484 ): 

485 self.log.warning( 

486 "Remote OCR was requested for this document but no remote " 

487 "parser is available for it, processing locally instead.", 

488 ) 

489 

490 # Notify all listeners that we're going to do some work. 

491 

492 document_consumption_started.send( 

493 sender=self.__class__, 

494 filename=self.working_copy, 

495 logging_group=self.logging_group, 

496 ) 

497 

498 self.run_pre_consume_script() 

499 

500 # This doesn't parse the document yet, but gives us a parser. 

501 with parser_class() as document_parser: 

502 document_parser.configure( 

503 ParserContext(mailrule_id=self.input_doc.mailrule_id), 

504 ) 

505 

506 self.log.debug( 

507 f"Parser: {document_parser.name} v{document_parser.version}", 

508 ) 

509 

510 # New versions skip the barcode plugin, so read their barcodes here 

511 if ( 

512 self.input_doc.root_document_id is not None 

513 and self.metadata.barcodes is None 

514 ): 

515 self._read_version_barcodes(mime_type, Path(tmpdir)) 

516 

517 # Parse the document. This may take some time. 

518 

519 text = None 

520 date = None 

521 thumbnail = None 

522 archive_path = None 

523 page_count = None 

524 

525 try: 

526 self._send_progress( 

527 20, 

528 100, 

529 ProgressStatusOptions.WORKING, 

530 ConsumerStatusShortMessage.PARSING_DOCUMENT, 

531 ) 

532 self.log.debug(f"Parsing {self.filename}...") 

533 

534 produce_archive = should_produce_archive( 

535 document_parser, 

536 mime_type, 

537 self.working_copy, 

538 self.log, 

539 ) 

540 document_parser.parse( 

541 self.working_copy, 

542 mime_type, 

543 produce_archive=produce_archive, 

544 ) 

545 

546 self.log.debug(f"Generating thumbnail for {self.filename}...") 

547 self._send_progress( 

548 70, 

549 100, 

550 ProgressStatusOptions.WORKING, 

551 ConsumerStatusShortMessage.GENERATING_THUMBNAIL, 

552 ) 

553 thumbnail = document_parser.get_thumbnail( 

554 self.working_copy, 

555 mime_type, 

556 ) 

557 

558 text = document_parser.get_text() 

559 date = document_parser.get_date() 

560 if date is None: 

561 self._send_progress( 

562 90, 

563 100, 

564 ProgressStatusOptions.WORKING, 

565 ConsumerStatusShortMessage.PARSE_DATE, 

566 ) 

567 with get_date_parser() as date_parser: 

568 date = next(date_parser.parse(self.filename, text), None) 

569 archive_path = document_parser.get_archive_path() 

570 page_count = document_parser.get_page_count( 

571 self.working_copy, 

572 mime_type, 

573 ) 

574 

575 except ParseError as e: 

576 self._fail( 

577 str(e), 

578 f"Error occurred while consuming document {self.filename}: {e}", 

579 exc_info=True, 

580 exception=e, 

581 ) 

582 except Exception as e: 

583 self._fail( 

584 str(e), 

585 f"Unexpected error while consuming document {self.filename}: {e}", 

586 exc_info=True, 

587 exception=e, 

588 ) 

589 

590 # Prepare the document classifier. 

591 

592 # TODO: I don't really like to do this here, but this way we avoid 

593 # reloading the classifier multiple times, since there are multiple 

594 # post-consume hooks that all require the classifier. 

595 

596 classifier = load_classifier() 

597 

598 self._send_progress( 

599 95, 

600 100, 

601 ProgressStatusOptions.WORKING, 

602 ConsumerStatusShortMessage.SAVE_DOCUMENT, 

603 ) 

604 # now that everything is done, we can start to store the document 

605 # in the system. This will be a transaction and reasonably fast. 

606 try: 

607 with transaction.atomic(): 

608 # store the document. 

609 if self.input_doc.root_document_id: 

610 # If this is a new version of an existing document, we need 

611 # to make sure we're not creating a new document, but updating 

612 # the existing one. 

613 root_doc = Document.objects.get( 

614 pk=self.input_doc.root_document_id, 

615 ) 

616 original_document = self._create_version_from_root( 

617 root_doc, 

618 text=text, 

619 page_count=page_count, 

620 mime_type=mime_type, 

621 ) 

622 actor = None 

623 

624 # Save the new version, potentially creating an audit log entry for the version addition if enabled. 

625 if ( 

626 settings.AUDIT_LOG_ENABLED 

627 and self.metadata.actor_id is not None 

628 ): 

629 actor = User.objects.filter( 

630 pk=self.metadata.actor_id, 

631 ).first() 

632 if actor is not None: 

633 from auditlog.context import ( # type: ignore[import-untyped] 

634 set_actor, 

635 ) 

636 

637 with set_actor(actor): 

638 original_document.save() 

639 else: 

640 original_document.save() 

641 else: 

642 original_document.save() 

643 

644 self._store_barcodes(original_document) 

645 

646 # Adding a version changes the effective document, so update root modified 

647 Document.objects.filter(pk=root_doc.pk).update( 

648 modified=timezone.now(), 

649 ) 

650 

651 # Create a log entry for the version addition, if enabled 

652 if settings.AUDIT_LOG_ENABLED: 

653 from auditlog.models import ( # type: ignore[import-untyped] 

654 LogEntry, 

655 ) 

656 

657 LogEntry.objects.log_create( 

658 instance=root_doc, 

659 changes={ 

660 "Version Added": ["None", original_document.id], 

661 }, 

662 action=LogEntry.Action.UPDATE, 

663 actor=actor, 

664 additional_data={ 

665 "reason": "Version added", 

666 "version_id": original_document.id, 

667 }, 

668 ) 

669 document = original_document 

670 else: 

671 document = self._store( 

672 text=text, 

673 date=date, 

674 page_count=page_count, 

675 mime_type=mime_type, 

676 ) 

677 

678 # If we get here, it was successful. Proceed with post-consume 

679 # hooks. If they fail, nothing will get changed. 

680 

681 document = Document.objects.prefetch_related("versions").get( 

682 pk=document.pk, 

683 ) 

684 

685 document_consumption_finished.send( 

686 sender=self.__class__, 

687 document=document, 

688 logging_group=self.logging_group, 

689 classifier=classifier, 

690 original_file=self.unmodified_original 

691 if self.unmodified_original 

692 else self.working_copy, 

693 ) 

694 

695 # After everything is in the database, copy the files into 

696 # place. If this fails, we'll also rollback the transaction. 

697 with FileLock(settings.MEDIA_LOCK): 

698 generated_filename = generate_unique_filename(document) 

699 if ( 

700 len(str(generated_filename)) 

701 > Document.MAX_STORED_FILENAME_LENGTH 

702 ): 

703 self.log.warning( 

704 "Generated source filename exceeds db path limit, falling back to default naming", 

705 ) 

706 generated_filename = generate_filename( 

707 document, 

708 use_format=False, 

709 ) 

710 document.filename = generated_filename 

711 validate_path_in_root( 

712 document.source_path, 

713 settings.ORIGINALS_DIR, 

714 ) 

715 create_source_path_directory(document.source_path) 

716 

717 self._write( 

718 self.unmodified_original 

719 if self.unmodified_original is not None 

720 else self.working_copy, 

721 document.source_path, 

722 ) 

723 

724 self._write( 

725 thumbnail, 

726 document.thumbnail_path, 

727 ) 

728 

729 if archive_path and Path(archive_path).is_file(): 

730 generated_archive_filename = generate_unique_filename( 

731 document, 

732 archive_filename=True, 

733 ) 

734 if ( 

735 len(str(generated_archive_filename)) 

736 > Document.MAX_STORED_FILENAME_LENGTH 

737 ): 

738 self.log.warning( 

739 "Generated archive filename exceeds db path limit, falling back to default naming", 

740 ) 

741 generated_archive_filename = generate_filename( 

742 document, 

743 archive_filename=True, 

744 use_format=False, 

745 ) 

746 document.archive_filename = generated_archive_filename 

747 validate_path_in_root( 

748 document.archive_path, 

749 settings.ARCHIVE_DIR, 

750 ) 

751 create_source_path_directory(document.archive_path) 

752 self._write( 

753 archive_path, 

754 document.archive_path, 

755 ) 

756 

757 document.archive_checksum = compute_checksum( 

758 document.archive_path, 

759 ) 

760 

761 # Don't save with the lock active. Saving will cause the file 

762 # renaming logic to acquire the lock as well. 

763 # This triggers things like file renaming 

764 document.save() 

765 

766 if document.root_document_id: 

767 document_updated.send( 

768 sender=self.__class__, 

769 document=document.root_document, 

770 skip_ai_index=True, # document_consumption_finished already enqueues the LLM update 

771 ) 

772 

773 # Delete the file only if it was successfully consumed 

774 self.log.debug( 

775 f"Deleting original file {self.input_doc.original_file}", 

776 ) 

777 self.input_doc.original_file.unlink() 

778 self.log.debug(f"Deleting working copy {self.working_copy}") 

779 self.working_copy.unlink() 

780 if self.unmodified_original is not None: # pragma: no cover 

781 self.log.debug( 

782 f"Deleting unmodified original file {self.unmodified_original}", 

783 ) 

784 self.unmodified_original.unlink() 

785 

786 # https://github.com/jonaswinkler/paperless-ng/discussions/1037 

787 shadow_file = ( 

788 Path(self.input_doc.original_file).parent 

789 / f"._{Path(self.input_doc.original_file).name}" 

790 ) 

791 

792 if Path(shadow_file).is_file(): 

793 self.log.debug(f"Deleting shadow file {shadow_file}") 

794 Path(shadow_file).unlink() 

795 

796 except Exception as e: 

797 self._fail( 

798 str(e), 

799 f"The following error occurred while storing document " 

800 f"{self.filename} after parsing: {e}", 

801 exc_info=True, 

802 exception=e, 

803 ) 

804 

805 self.run_post_consume_script(document) 

806 

807 self.log.info(f"Document {document} consumption finished") 

808 

809 self._send_progress( 

810 100, 

811 100, 

812 ProgressStatusOptions.SUCCESS, 

813 ConsumerStatusShortMessage.FINISHED, 

814 document.id, 

815 ) 

816 

817 # Return the most up to date fields 

818 document.refresh_from_db() 

819 

820 return ConsumeFileSuccessResult(document_id=document.pk) 

821 

822 def _parse_title_placeholders(self, title: str) -> str: 

823 local_added = timezone.localtime(timezone.now()) 

824 

825 correspondent_name = ( 

826 Correspondent.objects.get(pk=self.metadata.correspondent_id).name 

827 if self.metadata.correspondent_id is not None 

828 else None 

829 ) 

830 doc_type_name = ( 

831 DocumentType.objects.get(pk=self.metadata.document_type_id).name 

832 if self.metadata.document_type_id is not None 

833 else None 

834 ) 

835 owner_username = ( 

836 User.objects.get(pk=self.metadata.owner_id).username 

837 if self.metadata.owner_id is not None 

838 else None 

839 ) 

840 

841 return parse_w_workflow_placeholders( 

842 title, 

843 correspondent_name, 

844 doc_type_name, 

845 owner_username, 

846 local_added, 

847 self.filename, 

848 self.filename, 

849 ) 

850 

851 def _store( 

852 self, 

853 text: str, 

854 date: datetime.datetime | None, 

855 page_count: int | None, 

856 mime_type: str, 

857 ) -> Document: 

858 # If someone gave us the original filename, use it instead of doc. 

859 

860 self.log.debug("Saving record to database") 

861 

862 if self.metadata.created is not None: 

863 create_date = self.metadata.created 

864 self.log.debug( 

865 f"Creation date from post_documents parameter: {create_date}", 

866 ) 

867 elif date is not None: 

868 create_date = date 

869 self.log.debug(f"Creation date from parse_date: {create_date}") 

870 else: 

871 stats = Path(self.input_doc.original_file).stat() 

872 create_date = datetime.datetime.fromtimestamp( 

873 stats.st_mtime, 

874 tz=timezone.get_current_timezone(), 

875 ) 

876 self.log.debug(f"Creation date from st_mtime: {create_date}") 

877 

878 if self.metadata.filename: 

879 title = Path(self.metadata.filename).stem 

880 else: 

881 title = self.input_doc.original_file.stem 

882 

883 if self.metadata.title is not None: 

884 try: 

885 title = self._parse_title_placeholders(self.metadata.title) 

886 except Exception as e: 

887 self.log.error( 

888 f"Error occurred parsing title override '{self.metadata.title}', falling back to original. Exception: {e}", 

889 ) 

890 

891 file_for_checksum = ( 

892 self.unmodified_original 

893 if self.unmodified_original is not None 

894 else self.working_copy 

895 ) 

896 

897 document = Document.objects.create( 

898 title=title[:127], 

899 content=text, 

900 mime_type=mime_type, 

901 checksum=compute_checksum(file_for_checksum), 

902 created=create_date, 

903 modified=create_date, 

904 page_count=page_count, 

905 original_filename=self.filename, 

906 ) 

907 

908 self.apply_overrides(document) 

909 

910 document.save() 

911 

912 return document 

913 

914 def apply_overrides(self, document: Document) -> None: 

915 if self.metadata.correspondent_id: 

916 document.correspondent = Correspondent.objects.get( 

917 pk=self.metadata.correspondent_id, 

918 ) 

919 

920 if self.metadata.document_type_id: 

921 document.document_type = DocumentType.objects.get( 

922 pk=self.metadata.document_type_id, 

923 ) 

924 

925 if self.metadata.tag_ids: 

926 for tag_id in self.metadata.tag_ids: 

927 document.add_nested_tags([Tag.objects.get(pk=tag_id)]) 

928 

929 if self.metadata.storage_path_id: 

930 document.storage_path = StoragePath.objects.get( 

931 pk=self.metadata.storage_path_id, 

932 ) 

933 

934 if self.metadata.asn is not None: 

935 document.archive_serial_number = self.metadata.asn 

936 

937 if self.metadata.version_label is not None: 

938 document.version_label = self.metadata.version_label 

939 

940 if self.metadata.owner_id: 

941 document.owner = User.objects.get( 

942 pk=self.metadata.owner_id, 

943 ) 

944 

945 if ( 

946 self.metadata.view_users is not None 

947 or self.metadata.view_groups is not None 

948 or self.metadata.change_users is not None 

949 or self.metadata.change_groups is not None 

950 ): 

951 permissions = { 

952 "view": { 

953 "users": self.metadata.view_users or [], 

954 "groups": self.metadata.view_groups or [], 

955 }, 

956 "change": { 

957 "users": self.metadata.change_users or [], 

958 "groups": self.metadata.change_groups or [], 

959 }, 

960 } 

961 set_permissions_for_object(permissions=permissions, object=document) 

962 

963 if self.metadata.custom_fields: 

964 for field in CustomField.objects.filter( 

965 id__in=self.metadata.custom_fields.keys(), 

966 ).distinct(): 

967 value_field_name = CustomFieldInstance.get_value_field_name( 

968 data_type=field.data_type, 

969 ) 

970 args = { 

971 "field": field, 

972 "document": document, 

973 value_field_name: self.metadata.custom_fields.get(field.id, None), 

974 } 

975 CustomFieldInstance.objects.create(**args) # adds to document 

976 

977 self._store_barcodes(document) 

978 

979 def _read_version_barcodes(self, mime_type: str, work_dir: Path) -> None: 

980 barcode_settings = BarcodeConfig() 

981 if not barcode_settings.barcode_store_values: 

982 return 

983 try: 

984 self.metadata.barcodes = ( 

985 read_barcode_values( 

986 self.working_copy, 

987 mime_type, 

988 barcode_settings, 

989 work_dir, 

990 ) 

991 or None 

992 ) 

993 except Exception as e: 

994 self.log.warning(f"Could not read barcodes of {self.filename}: {e}") 

995 

996 def _store_barcodes(self, document: Document) -> None: 

997 if self.metadata.barcodes: 

998 DocumentBarcode.objects.bulk_create( 

999 DocumentBarcode(document=document, **barcode) 

1000 for barcode in self.metadata.barcodes 

1001 ) 

1002 

1003 def _write(self, source, target) -> None: 

1004 with ( 

1005 Path(source).open("rb") as read_file, 

1006 Path(target).open("wb") as write_file, 

1007 ): 

1008 shutil.copyfileobj(read_file, write_file) 

1009 

1010 # Attempt to copy file's original stats, but it's ok if we can't 

1011 try: 

1012 copy_basic_file_stats(source, target) 

1013 except Exception: # pragma: no cover 

1014 pass 

1015 

1016 

1017class ConsumerPreflightPlugin( 

1018 NoCleanupPluginMixin, 

1019 NoSetupPluginMixin, 

1020 AlwaysRunPluginMixin, 

1021 LoggingMixin, 

1022 ConsumerPluginMixin, 

1023 ConsumeTaskPlugin, 

1024): 

1025 NAME: str = "ConsumerPreflightPlugin" 

1026 logging_name = LOGGING_NAME 

1027 

1028 def pre_check_file_exists(self) -> None: 

1029 """ 

1030 Confirm the input file still exists where it should 

1031 """ 

1032 if TYPE_CHECKING: 

1033 assert isinstance(self.input_doc.original_file, Path), ( 

1034 self.input_doc.original_file 

1035 ) 

1036 if not self.input_doc.original_file.is_file(): 

1037 self._fail( 

1038 ConsumerStatusShortMessage.FILE_NOT_FOUND, 

1039 f"Cannot consume {self.input_doc.original_file}: File not found.", 

1040 ) 

1041 

1042 def pre_check_duplicate(self) -> None: 

1043 """ 

1044 Using the SHA256 of the file, check this exact file doesn't already exist 

1045 """ 

1046 checksum = compute_checksum(Path(self.input_doc.original_file)) 

1047 existing_doc = Document.global_objects.filter( 

1048 Q(checksum=checksum) | Q(archive_checksum=checksum), 

1049 ) 

1050 if existing_doc.exists(): 

1051 existing_doc = existing_doc.order_by("-created") 

1052 duplicates_in_trash = existing_doc.filter(deleted_at__isnull=False) 

1053 log_msg = ( 

1054 f"Consuming duplicate {self.filename}: " 

1055 f"{existing_doc.count()} existing document(s) share the same content." 

1056 ) 

1057 

1058 if duplicates_in_trash.exists(): 

1059 log_msg += " Note: at least one existing document is in the trash." 

1060 

1061 self.log.warning(log_msg) 

1062 

1063 if settings.CONSUMER_DELETE_DUPLICATES: 

1064 duplicate = existing_doc.first() 

1065 duplicate_label = ( 

1066 duplicate.title 

1067 or duplicate.original_filename 

1068 or (Path(duplicate.filename).name if duplicate.filename else None) 

1069 or str(duplicate.pk) 

1070 ) 

1071 

1072 Path(self.input_doc.original_file).unlink() 

1073 

1074 failure_msg = ( 

1075 f"Not consuming {self.filename}: " 

1076 f"It is a duplicate of {duplicate_label} (#{duplicate.pk})" 

1077 ) 

1078 status_msg = ConsumerStatusShortMessage.DOCUMENT_ALREADY_EXISTS 

1079 

1080 if duplicates_in_trash.exists(): 

1081 status_msg = ( 

1082 ConsumerStatusShortMessage.DOCUMENT_ALREADY_EXISTS_IN_TRASH 

1083 ) 

1084 failure_msg += " Note: existing document is in the trash." 

1085 

1086 self._send_progress(100, 100, ProgressStatusOptions.FAILED, status_msg) 

1087 self.log.error(failure_msg) 

1088 in_trash = duplicates_in_trash.exists() 

1089 raise ConsumeFileDuplicateError( 

1090 f"{self.filename}: {failure_msg}", 

1091 duplicate.pk, 

1092 in_trash=in_trash, 

1093 ) 

1094 

1095 def pre_check_directories(self) -> None: 

1096 """ 

1097 Ensure all required directories exist before attempting to use them 

1098 """ 

1099 settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True) 

1100 settings.THUMBNAIL_DIR.mkdir(parents=True, exist_ok=True) 

1101 settings.ORIGINALS_DIR.mkdir(parents=True, exist_ok=True) 

1102 settings.ARCHIVE_DIR.mkdir(parents=True, exist_ok=True) 

1103 

1104 def run(self) -> None: 

1105 self._send_progress( 

1106 0, 

1107 100, 

1108 ProgressStatusOptions.STARTED, 

1109 ConsumerStatusShortMessage.NEW_FILE, 

1110 ) 

1111 

1112 # Make sure that preconditions for consuming the file are met. 

1113 

1114 self.pre_check_file_exists() 

1115 self.pre_check_duplicate() 

1116 self.pre_check_directories() 

1117 

1118 

1119class AsnCheckPlugin( 

1120 NoCleanupPluginMixin, 

1121 NoSetupPluginMixin, 

1122 AlwaysRunPluginMixin, 

1123 LoggingMixin, 

1124 ConsumerPluginMixin, 

1125 ConsumeTaskPlugin, 

1126): 

1127 NAME: str = "AsnCheckPlugin" 

1128 logging_name = LOGGING_NAME 

1129 

1130 def pre_check_asn_value(self) -> None: 

1131 """ 

1132 Check that if override_asn is given, it is unique and within a valid range 

1133 """ 

1134 if self.metadata.asn is None: 

1135 # if ASN is None 

1136 return 

1137 # Validate the range is above zero and less than uint32_t max 

1138 # otherwise, Whoosh can't handle it in the index 

1139 if ( 

1140 self.metadata.asn < Document.ARCHIVE_SERIAL_NUMBER_MIN 

1141 or self.metadata.asn > Document.ARCHIVE_SERIAL_NUMBER_MAX 

1142 ): 

1143 self._fail( 

1144 ConsumerStatusShortMessage.ASN_RANGE, 

1145 f"Not consuming {self.filename}: " 

1146 f"Given ASN {self.metadata.asn} is out of range " 

1147 f"[{Document.ARCHIVE_SERIAL_NUMBER_MIN:,}, " 

1148 f"{Document.ARCHIVE_SERIAL_NUMBER_MAX:,}]", 

1149 ) 

1150 existing_asn_doc = Document.global_objects.filter( 

1151 archive_serial_number=self.metadata.asn, 

1152 ) 

1153 if existing_asn_doc.exists(): 

1154 msg = ConsumerStatusShortMessage.ASN_ALREADY_EXISTS 

1155 log_msg = f"Not consuming {self.filename}: Given ASN {self.metadata.asn} already exists!" 

1156 

1157 if existing_asn_doc.first().deleted_at is not None: 

1158 msg = ConsumerStatusShortMessage.ASN_ALREADY_EXISTS_IN_TRASH 

1159 log_msg += " Note: existing document is in the trash." 

1160 

1161 self._fail( 

1162 msg, 

1163 log_msg, 

1164 ) 

1165 

1166 def run(self) -> None: 

1167 self.pre_check_asn_value()