Coverage for documents/consumer.py: 21%
449 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
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
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
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
67LOGGING_NAME: Final[str] = "paperless.consumer"
70class WorkflowTriggerPlugin(
71 NoCleanupPluginMixin,
72 NoSetupPluginMixin,
73 AlwaysRunPluginMixin,
74 ConsumeTaskPlugin,
75):
76 NAME: str = "WorkflowTriggerPlugin"
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
93class ConsumerError(Exception):
94 pass
97class ConsumeFileDuplicateError(ConsumerError):
98 """Raised when a file is rejected because it duplicates an existing document."""
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
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"
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.
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)
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
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
152 generation = OcrConfig().archive_file_generation
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
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
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
193 log: "LoggerAdapter" # type: ignore[type-arg]
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)
205 self.renew_logging_group()
207 self.filename = self.metadata.filename or self.input_doc.original_file.name
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 )
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
242class ConsumerPlugin(
243 AlwaysRunPluginMixin,
244 NoSetupPluginMixin,
245 NoCleanupPluginMixin,
246 LoggingMixin,
247 ConsumerPluginMixin,
248 ConsumeTaskPlugin,
249):
250 logging_name = LOGGING_NAME
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
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
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 )
308 self.log.info(f"Executing pre-consume script {settings.PRE_CONSUME_SCRIPT}")
310 working_file_path = str(self.working_copy)
311 original_file_path = str(self.input_doc.original_file)
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 ""
318 try:
319 run_subprocess(
320 [
321 settings.PRE_CONSUME_SCRIPT,
322 ],
323 script_env,
324 self.log,
325 )
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 )
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
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 )
350 self.log.info(
351 f"Executing post-consume script {settings.POST_CONSUME_SCRIPT}",
352 )
354 script_env = os.environ.copy()
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 ""
387 try:
388 run_subprocess(
389 [
390 settings.POST_CONSUME_SCRIPT,
391 ],
392 script_env,
393 self.log,
394 )
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 )
404 def run(self) -> "ConsumeFileSuccessResult":
405 """
406 Return the document object if it was successfully created.
407 """
409 # Preflight has already run including progress update to 0%
410 self.log.info(f"Consuming {self.filename}")
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
421 # Determine the parser class.
423 mime_type = magic.from_file(self.working_copy, mime=True)
425 self.log.debug(f"Detected mime type: {mime_type}")
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}")
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 )
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 )
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 )
490 # Notify all listeners that we're going to do some work.
492 document_consumption_started.send(
493 sender=self.__class__,
494 filename=self.working_copy,
495 logging_group=self.logging_group,
496 )
498 self.run_pre_consume_script()
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 )
506 self.log.debug(
507 f"Parser: {document_parser.name} v{document_parser.version}",
508 )
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))
517 # Parse the document. This may take some time.
519 text = None
520 date = None
521 thumbnail = None
522 archive_path = None
523 page_count = None
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}...")
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 )
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 )
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 )
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 )
590 # Prepare the document classifier.
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.
596 classifier = load_classifier()
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
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 )
637 with set_actor(actor):
638 original_document.save()
639 else:
640 original_document.save()
641 else:
642 original_document.save()
644 self._store_barcodes(original_document)
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 )
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 )
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 )
678 # If we get here, it was successful. Proceed with post-consume
679 # hooks. If they fail, nothing will get changed.
681 document = Document.objects.prefetch_related("versions").get(
682 pk=document.pk,
683 )
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 )
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)
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 )
724 self._write(
725 thumbnail,
726 document.thumbnail_path,
727 )
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 )
757 document.archive_checksum = compute_checksum(
758 document.archive_path,
759 )
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()
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 )
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()
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 )
792 if Path(shadow_file).is_file():
793 self.log.debug(f"Deleting shadow file {shadow_file}")
794 Path(shadow_file).unlink()
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 )
805 self.run_post_consume_script(document)
807 self.log.info(f"Document {document} consumption finished")
809 self._send_progress(
810 100,
811 100,
812 ProgressStatusOptions.SUCCESS,
813 ConsumerStatusShortMessage.FINISHED,
814 document.id,
815 )
817 # Return the most up to date fields
818 document.refresh_from_db()
820 return ConsumeFileSuccessResult(document_id=document.pk)
822 def _parse_title_placeholders(self, title: str) -> str:
823 local_added = timezone.localtime(timezone.now())
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 )
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 )
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.
860 self.log.debug("Saving record to database")
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}")
878 if self.metadata.filename:
879 title = Path(self.metadata.filename).stem
880 else:
881 title = self.input_doc.original_file.stem
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 )
891 file_for_checksum = (
892 self.unmodified_original
893 if self.unmodified_original is not None
894 else self.working_copy
895 )
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 )
908 self.apply_overrides(document)
910 document.save()
912 return document
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 )
920 if self.metadata.document_type_id:
921 document.document_type = DocumentType.objects.get(
922 pk=self.metadata.document_type_id,
923 )
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)])
929 if self.metadata.storage_path_id:
930 document.storage_path = StoragePath.objects.get(
931 pk=self.metadata.storage_path_id,
932 )
934 if self.metadata.asn is not None:
935 document.archive_serial_number = self.metadata.asn
937 if self.metadata.version_label is not None:
938 document.version_label = self.metadata.version_label
940 if self.metadata.owner_id:
941 document.owner = User.objects.get(
942 pk=self.metadata.owner_id,
943 )
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)
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
977 self._store_barcodes(document)
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}")
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 )
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)
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
1017class ConsumerPreflightPlugin(
1018 NoCleanupPluginMixin,
1019 NoSetupPluginMixin,
1020 AlwaysRunPluginMixin,
1021 LoggingMixin,
1022 ConsumerPluginMixin,
1023 ConsumeTaskPlugin,
1024):
1025 NAME: str = "ConsumerPreflightPlugin"
1026 logging_name = LOGGING_NAME
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 )
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 )
1058 if duplicates_in_trash.exists():
1059 log_msg += " Note: at least one existing document is in the trash."
1061 self.log.warning(log_msg)
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 )
1072 Path(self.input_doc.original_file).unlink()
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
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."
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 )
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)
1104 def run(self) -> None:
1105 self._send_progress(
1106 0,
1107 100,
1108 ProgressStatusOptions.STARTED,
1109 ConsumerStatusShortMessage.NEW_FILE,
1110 )
1112 # Make sure that preconditions for consuming the file are met.
1114 self.pre_check_file_exists()
1115 self.pre_check_duplicate()
1116 self.pre_check_directories()
1119class AsnCheckPlugin(
1120 NoCleanupPluginMixin,
1121 NoSetupPluginMixin,
1122 AlwaysRunPluginMixin,
1123 LoggingMixin,
1124 ConsumerPluginMixin,
1125 ConsumeTaskPlugin,
1126):
1127 NAME: str = "AsnCheckPlugin"
1128 logging_name = LOGGING_NAME
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!"
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."
1161 self._fail(
1162 msg,
1163 log_msg,
1164 )
1166 def run(self) -> None:
1167 self.pre_check_asn_value()