Coverage for documents/signals/handlers.py: 37%
608 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 datetime
4import logging
5import shutil
6import traceback as _tb
7from pathlib import Path
8from typing import TYPE_CHECKING
9from typing import Any
11from celery import shared_task
12from celery.signals import before_task_publish
13from celery.signals import task_failure
14from celery.signals import task_postrun
15from celery.signals import task_prerun
16from celery.signals import task_revoked
17from celery.signals import worker_process_init
18from celery.signals import worker_process_shutdown
19from django.conf import settings
20from django.contrib.auth.models import Group
21from django.contrib.auth.models import User
22from django.db import DatabaseError
23from django.db import close_old_connections
24from django.db import connections
25from django.db import models
26from django.db.models import Q
27from django.dispatch import receiver
28from django.utils import timezone
29from filelock import FileLock
30from rest_framework import serializers
32from documents import matching
33from documents.caching import clear_document_caches
34from documents.caching import invalidate_llm_suggestions_cache
35from documents.caching import invalidate_suggestions_cache
36from documents.data_models import ConsumableDocument
37from documents.file_handling import create_source_path_directory
38from documents.file_handling import delete_empty_directories
39from documents.file_handling import generate_filename
40from documents.file_handling import generate_unique_filename
41from documents.models import Correspondent
42from documents.models import CustomField
43from documents.models import CustomFieldInstance
44from documents.models import Document
45from documents.models import DocumentType
46from documents.models import PaperlessTask
47from documents.models import SavedView
48from documents.models import StoragePath
49from documents.models import Tag
50from documents.models import UiSettings
51from documents.models import Workflow
52from documents.models import WorkflowAction
53from documents.models import WorkflowRun
54from documents.models import WorkflowTrigger
55from documents.permissions import get_objects_for_user_owner_aware
56from documents.plugins.helpers import DocumentsStatusManager
57from documents.templating.utils import convert_format_str_to_template_format
58from documents.utils import compute_checksum
59from documents.utils import copy_file_with_basic_stats
60from documents.workflows.actions import build_workflow_action_context
61from documents.workflows.actions import execute_email_action
62from documents.workflows.actions import execute_move_to_trash_action
63from documents.workflows.actions import execute_password_removal_action
64from documents.workflows.actions import execute_webhook_action
65from documents.workflows.mutations import apply_assignment_to_document
66from documents.workflows.mutations import apply_assignment_to_overrides
67from documents.workflows.mutations import apply_removal_to_document
68from documents.workflows.mutations import apply_removal_to_overrides
69from documents.workflows.utils import get_workflows_for_trigger
70from paperless.config import AIConfig
72if TYPE_CHECKING: 72 ↛ 73line 72 didn't jump to line 73 because the condition on line 72 was never true
73 import uuid
75 from documents.classifier import DocumentClassifier
76 from documents.data_models import ConsumableDocument
77 from documents.data_models import DocumentMetadataOverrides
79logger = logging.getLogger("paperless.handlers")
80DRF_DATETIME_FIELD = serializers.DateTimeField()
83def add_inbox_tags(sender, document: Document, logging_group=None, **kwargs) -> None:
84 if document.owner is not None:
85 tags = get_objects_for_user_owner_aware(
86 document.owner,
87 "documents.view_tag",
88 Tag,
89 )
90 else:
91 tags = Tag.objects.all()
92 inbox_tags = tags.filter(is_inbox_tag=True)
93 document.add_nested_tags(inbox_tags)
96def set_correspondent(
97 sender: object,
98 document: Document,
99 *,
100 logging_group: object = None,
101 classifier: DocumentClassifier | None = None,
102 replace: bool = False,
103 use_first: bool = True,
104 dry_run: bool = False,
105 **kwargs: Any,
106) -> Correspondent | None:
107 """
108 Assign a correspondent to a document based on classifier results.
110 Args:
111 document: The document to classify.
112 logging_group: Optional logging group for structured log output.
113 classifier: The trained classifier. If None, only rule-based matching runs.
114 replace: If True, overwrite an existing correspondent assignment.
115 use_first: If True, pick the first match when multiple correspondents
116 match. If False, skip assignment when multiple match.
117 dry_run: If True, compute and return the selection without saving.
118 **kwargs: Absorbed for Django signal compatibility (e.g. sender, signal).
120 Returns:
121 The correspondent that was (or would be) assigned, or None if no match
122 was found or assignment was skipped.
123 """
124 if document.correspondent and not replace:
125 return None
127 potential_correspondents = matching.match_correspondents(document, classifier)
128 potential_count = len(potential_correspondents)
129 selected = potential_correspondents[0] if potential_correspondents else None
131 if potential_count > 1:
132 if use_first:
133 logger.debug(
134 f"Detected {potential_count} potential correspondents, "
135 f"so we've opted for {selected}",
136 extra={"group": logging_group},
137 )
138 else:
139 logger.debug(
140 f"Detected {potential_count} potential correspondents, "
141 f"not assigning any correspondent",
142 extra={"group": logging_group},
143 )
144 return None
146 if (selected or replace) and not dry_run:
147 logger.info(
148 f"Assigning correspondent {selected} to {document}",
149 extra={"group": logging_group},
150 )
151 document.correspondent = selected
152 document.save(update_fields=("correspondent",))
154 return selected
157def set_document_type(
158 sender: object,
159 document: Document,
160 *,
161 logging_group: object = None,
162 classifier: DocumentClassifier | None = None,
163 replace: bool = False,
164 use_first: bool = True,
165 dry_run: bool = False,
166 **kwargs: Any,
167) -> DocumentType | None:
168 """
169 Assign a document type to a document based on classifier results.
171 Args:
172 document: The document to classify.
173 logging_group: Optional logging group for structured log output.
174 classifier: The trained classifier. If None, only rule-based matching runs.
175 replace: If True, overwrite an existing document type assignment.
176 use_first: If True, pick the first match when multiple types match.
177 If False, skip assignment when multiple match.
178 dry_run: If True, compute and return the selection without saving.
179 **kwargs: Absorbed for Django signal compatibility (e.g. sender, signal).
181 Returns:
182 The document type that was (or would be) assigned, or None if no match
183 was found or assignment was skipped.
184 """
185 if document.document_type and not replace:
186 return None
188 potential_document_types = matching.match_document_types(document, classifier)
189 potential_count = len(potential_document_types)
190 selected = potential_document_types[0] if potential_document_types else None
192 if potential_count > 1:
193 if use_first:
194 logger.info(
195 f"Detected {potential_count} potential document types, "
196 f"so we've opted for {selected}",
197 extra={"group": logging_group},
198 )
199 else:
200 logger.info(
201 f"Detected {potential_count} potential document types, "
202 f"not assigning any document type",
203 extra={"group": logging_group},
204 )
205 return None
207 if (selected or replace) and not dry_run:
208 logger.info(
209 f"Assigning document type {selected} to {document}",
210 extra={"group": logging_group},
211 )
212 document.document_type = selected
213 document.save(update_fields=("document_type",))
215 return selected
218def set_tags(
219 sender: object,
220 document: Document,
221 *,
222 logging_group: object = None,
223 classifier: DocumentClassifier | None = None,
224 replace: bool = False,
225 dry_run: bool = False,
226 **kwargs: Any,
227) -> tuple[set[Tag], set[Tag]]:
228 """
229 Assign tags to a document based on classifier results.
231 When replace=True, existing auto-matched and rule-matched tags are removed
232 before applying the new set (inbox tags and manually-added tags are preserved).
234 Args:
235 document: The document to classify.
236 logging_group: Optional logging group for structured log output.
237 classifier: The trained classifier. If None, only rule-based matching runs.
238 replace: If True, remove existing classifier-managed tags before applying
239 new ones. Inbox tags and manually-added tags are always preserved.
240 dry_run: If True, compute what would change without saving anything.
241 **kwargs: Absorbed for Django signal compatibility (e.g. sender, signal).
243 Returns:
244 A two-tuple of (tags_added, tags_removed). In non-replace mode,
245 tags_removed is always an empty set. In dry_run mode, neither set
246 is applied to the database.
247 """
248 # Compute which tags would be removed under replace mode.
249 # The filter mirrors the .delete() call below: keep inbox tags and
250 # manually-added tags (match="" and not auto-matched).
251 if replace:
252 tags_to_remove: set[Tag] = set(
253 document.tags.exclude(
254 is_inbox_tag=True,
255 ).exclude(
256 Q(match="") & ~Q(matching_algorithm=Tag.MATCH_AUTO),
257 ),
258 )
259 else:
260 tags_to_remove = set()
262 if replace and not dry_run:
263 Document.tags.through.objects.filter(document=document).exclude(
264 Q(tag__is_inbox_tag=True),
265 ).exclude(
266 Q(tag__match="") & ~Q(tag__matching_algorithm=Tag.MATCH_AUTO),
267 ).delete()
269 current_tags = set(document.tags.all())
270 matched_tags = matching.match_tags(document, classifier)
271 tags_to_add = set(matched_tags) - current_tags
273 if tags_to_add and not dry_run:
274 logger.info(
275 f'Tagging "{document}" with "{", ".join(t.name for t in tags_to_add)}"',
276 extra={"group": logging_group},
277 )
278 document.add_nested_tags(tags_to_add)
280 return tags_to_add, tags_to_remove
283def set_storage_path(
284 sender: object,
285 document: Document,
286 *,
287 logging_group: object = None,
288 classifier: DocumentClassifier | None = None,
289 replace: bool = False,
290 use_first: bool = True,
291 dry_run: bool = False,
292 **kwargs: Any,
293) -> StoragePath | None:
294 """
295 Assign a storage path to a document based on classifier results.
297 Args:
298 document: The document to classify.
299 logging_group: Optional logging group for structured log output.
300 classifier: The trained classifier. If None, only rule-based matching runs.
301 replace: If True, overwrite an existing storage path assignment.
302 use_first: If True, pick the first match when multiple paths match.
303 If False, skip assignment when multiple match.
304 dry_run: If True, compute and return the selection without saving.
305 **kwargs: Absorbed for Django signal compatibility (e.g. sender, signal).
307 Returns:
308 The storage path that was (or would be) assigned, or None if no match
309 was found or assignment was skipped.
310 """
311 if document.storage_path and not replace:
312 return None
314 potential_storage_paths = matching.match_storage_paths(document, classifier)
315 potential_count = len(potential_storage_paths)
316 selected = potential_storage_paths[0] if potential_storage_paths else None
318 if potential_count > 1:
319 if use_first:
320 logger.info(
321 f"Detected {potential_count} potential storage paths, "
322 f"so we've opted for {selected}",
323 extra={"group": logging_group},
324 )
325 else:
326 logger.info(
327 f"Detected {potential_count} potential storage paths, "
328 f"not assigning any storage directory",
329 extra={"group": logging_group},
330 )
331 return None
333 if (selected or replace) and not dry_run:
334 logger.info(
335 f"Assigning storage path {selected} to {document}",
336 extra={"group": logging_group},
337 )
338 document.storage_path = selected
339 document.save(update_fields=("storage_path",))
341 return selected
344# see empty_trash in documents/tasks.py for signal handling
345def cleanup_document_deletion(sender, instance, **kwargs) -> None:
346 with FileLock(settings.MEDIA_LOCK):
347 if settings.EMPTY_TRASH_DIR: 347 ↛ 350line 347 didn't jump to line 350 because the condition on line 347 was never true
348 # Find a non-conflicting filename in case a document with the same
349 # name was moved to trash earlier
350 counter = 0
351 old_filename = Path(instance.source_path).name
352 old_filebase = Path(old_filename).stem
353 old_fileext = Path(old_filename).suffix
355 while True:
356 new_file_path = settings.EMPTY_TRASH_DIR / (
357 old_filebase + (f"_{counter:02}" if counter else "") + old_fileext
358 )
360 if new_file_path.exists():
361 counter += 1
362 else:
363 break
365 logger.debug(f"Moving {instance.source_path} to trash at {new_file_path}")
366 try:
367 shutil.move(
368 instance.source_path,
369 new_file_path,
370 copy_function=copy_file_with_basic_stats,
371 )
372 except OSError as e:
373 logger.error(
374 f"Failed to move {instance.source_path} to trash at "
375 f"{new_file_path}: {e}. Skipping cleanup!",
376 )
377 return
379 files = (
380 instance.archive_path,
381 instance.thumbnail_path,
382 )
383 if not settings.EMPTY_TRASH_DIR: 383 ↛ 387line 383 didn't jump to line 387 because the condition on line 383 was always true
384 # Only delete the original file if we are not moving it to trash dir
385 files += (instance.source_path,)
387 for filename in files:
388 if filename and filename.is_file():
389 try:
390 filename.unlink()
391 logger.debug(f"Deleted file {filename}.")
392 except OSError as e:
393 logger.warning(
394 f"While deleting document {instance!s}, the file "
395 f"{filename} could not be deleted: {e}",
396 )
397 elif filename and not filename.is_file():
398 logger.warning(f"Expected {filename} to exist, but it did not")
400 delete_empty_directories(
401 Path(instance.source_path).parent,
402 root=settings.ORIGINALS_DIR,
403 )
405 if instance.has_archive_version: 405 ↛ 406line 405 didn't jump to line 406 because the condition on line 405 was never true
406 delete_empty_directories(
407 Path(instance.archive_path).parent,
408 root=settings.ARCHIVE_DIR,
409 )
412class CannotMoveFilesException(Exception):
413 pass
416def _path_matches_checksum(path: Path, checksum: str | None) -> bool:
417 if checksum is None or not path.is_file():
418 return False
420 return compute_checksum(path) == checksum
423def _filename_template_uses_custom_fields(doc: Document) -> bool:
424 template = None
425 if doc.storage_path is not None:
426 template = doc.storage_path.path
427 elif settings.FILENAME_FORMAT is not None:
428 template = convert_format_str_to_template_format(settings.FILENAME_FORMAT)
430 if not template:
431 return False
433 return "custom_fields" in template
436# should be disabled in /src/documents/management/commands/document_importer.py handle
437@receiver(models.signals.post_save, sender=CustomFieldInstance, weak=False)
438@receiver(models.signals.m2m_changed, sender=Document.tags.through, weak=False)
439@receiver(models.signals.post_save, sender=Document, weak=False)
440def update_filename_and_move_files(
441 sender,
442 instance: Document | CustomFieldInstance,
443 **kwargs,
444) -> None:
445 if isinstance(instance, CustomFieldInstance): 445 ↛ 446line 445 didn't jump to line 446 because the condition on line 445 was never true
446 if not _filename_template_uses_custom_fields(instance.document):
447 return
448 instance = instance.document
450 def validate_move(instance, old_path: Path, new_path: Path, root: Path) -> None:
451 if not new_path.is_relative_to(root): 451 ↛ 452line 451 didn't jump to line 452 because the condition on line 451 was never true
452 msg = (
453 f"Document {instance!s}: Refusing to move file outside root {root}: "
454 f"{new_path}."
455 )
456 logger.warning(msg)
457 raise CannotMoveFilesException(msg)
459 if not old_path.is_file(): 459 ↛ 461line 459 didn't jump to line 461 because the condition on line 459 was never true
460 # Can't do anything if the old file does not exist anymore.
461 msg = f"Document {instance!s}: File {old_path} doesn't exist."
462 logger.fatal(msg)
463 raise CannotMoveFilesException(msg)
465 if new_path.is_file(): 465 ↛ 467line 465 didn't jump to line 467 because the condition on line 465 was never true
466 # Can't do anything if the new file already exists. Skip updating file.
467 msg = f"Document {instance!s}: Cannot rename file since target path {new_path} already exists."
468 logger.warning(msg)
469 raise CannotMoveFilesException(msg)
471 if not instance.filename: 471 ↛ 480line 471 didn't jump to line 480 because the condition on line 471 was never true
472 # Can't update the filename if there is no filename to begin with
473 # This happens when the consumer creates a new document.
474 # The document is modified and saved multiple times, and only after
475 # everything is done (i.e., the generated filename is final),
476 # filename will be set to the location where the consumer has put
477 # the file.
478 #
479 # This will in turn cause this logic to move the file where it belongs.
480 return
482 with FileLock(settings.MEDIA_LOCK):
483 try:
484 # If this was waiting for the lock, the filename or archive_filename
485 # of this document may have been updated. This happens if multiple updates
486 # get queued from the UI for the same document
487 # So freshen up the data before doing anything
488 instance.refresh_from_db()
490 old_filename = instance.filename
491 old_source_path = instance.source_path
492 move_original = False
493 original_already_moved = False
495 old_archive_filename = instance.archive_filename
496 old_archive_path = instance.archive_path
497 move_archive = False
498 archive_already_moved = False
500 candidate_filename = generate_filename(instance)
501 if len(str(candidate_filename)) > Document.MAX_STORED_FILENAME_LENGTH: 501 ↛ 502line 501 didn't jump to line 502 because the condition on line 501 was never true
502 msg = (
503 f"Document {instance!s}: Generated filename exceeds db path "
504 f"limit ({len(str(candidate_filename))} > "
505 f"{Document.MAX_STORED_FILENAME_LENGTH}): {candidate_filename!s}"
506 )
507 logger.warning(msg)
508 raise CannotMoveFilesException(msg)
510 candidate_source_path = (
511 settings.ORIGINALS_DIR / candidate_filename
512 ).resolve()
513 if candidate_filename == Path(old_filename):
514 new_filename = Path(old_filename)
515 elif ( 515 ↛ 519line 515 didn't jump to line 519 because the condition on line 515 was never true
516 candidate_source_path.exists()
517 and candidate_source_path != old_source_path
518 ):
519 if not old_source_path.is_file() and _path_matches_checksum(
520 candidate_source_path,
521 instance.checksum,
522 ):
523 new_filename = candidate_filename
524 original_already_moved = True
525 else:
526 # Only fall back to unique search when there is an actual conflict
527 new_filename = generate_unique_filename(instance)
528 else:
529 new_filename = candidate_filename
531 # Need to convert to string to be able to save it to the db
532 instance.filename = str(new_filename)
533 move_original = (
534 old_filename != instance.filename and not original_already_moved
535 )
537 if instance.has_archive_version: 537 ↛ 538line 537 didn't jump to line 538 because the condition on line 537 was never true
538 archive_candidate = generate_filename(instance, archive_filename=True)
539 if len(str(archive_candidate)) > Document.MAX_STORED_FILENAME_LENGTH:
540 msg = (
541 f"Document {instance!s}: Generated archive filename exceeds "
542 f"db path limit ({len(str(archive_candidate))} > "
543 f"{Document.MAX_STORED_FILENAME_LENGTH}): {archive_candidate!s}"
544 )
545 logger.warning(msg)
546 raise CannotMoveFilesException(msg)
547 archive_candidate_path = (
548 settings.ARCHIVE_DIR / archive_candidate
549 ).resolve()
550 if archive_candidate == Path(old_archive_filename):
551 new_archive_filename = Path(old_archive_filename)
552 elif (
553 archive_candidate_path.exists()
554 and archive_candidate_path != old_archive_path
555 ):
556 if not old_archive_path.is_file() and _path_matches_checksum(
557 archive_candidate_path,
558 instance.archive_checksum,
559 ):
560 new_archive_filename = archive_candidate
561 archive_already_moved = True
562 else:
563 new_archive_filename = generate_unique_filename(
564 instance,
565 archive_filename=True,
566 )
567 else:
568 new_archive_filename = archive_candidate
570 instance.archive_filename = str(new_archive_filename)
572 move_archive = (
573 old_archive_filename != instance.archive_filename
574 and not archive_already_moved
575 )
576 else:
577 move_archive = False
579 if not move_original and not move_archive:
580 updates = {"modified": timezone.now()}
581 if old_filename != instance.filename: 581 ↛ 582line 581 didn't jump to line 582 because the condition on line 581 was never true
582 updates["filename"] = instance.filename
583 if old_archive_filename != instance.archive_filename: 583 ↛ 584line 583 didn't jump to line 584 because the condition on line 583 was never true
584 updates["archive_filename"] = instance.archive_filename
586 # Don't save() here to prevent infinite recursion.
587 Document.objects.filter(pk=instance.pk).update(**updates)
588 return
590 if move_original: 590 ↛ 600line 590 didn't jump to line 600 because the condition on line 590 was always true
591 validate_move(
592 instance,
593 old_source_path,
594 instance.source_path,
595 settings.ORIGINALS_DIR,
596 )
597 create_source_path_directory(instance.source_path)
598 shutil.move(old_source_path, instance.source_path)
600 if move_archive: 600 ↛ 601line 600 didn't jump to line 601 because the condition on line 600 was never true
601 validate_move(
602 instance,
603 old_archive_path,
604 instance.archive_path,
605 settings.ARCHIVE_DIR,
606 )
607 create_source_path_directory(instance.archive_path)
608 shutil.move(old_archive_path, instance.archive_path)
610 # Don't save() here to prevent infinite recursion.
611 Document.global_objects.filter(pk=instance.pk).update(
612 filename=instance.filename,
613 archive_filename=instance.archive_filename,
614 modified=timezone.now(),
615 )
616 # Clear any caching for this document. Slightly overkill, but not terrible
617 clear_document_caches(instance.pk)
619 except (OSError, DatabaseError, CannotMoveFilesException) as e:
620 logger.warning(f"Exception during file handling: {e}")
621 # This happens when either:
622 # - moving the files failed due to file system errors
623 # - saving to the database failed due to database errors
624 # In both cases, we need to revert to the original state.
626 # Try to move files to their original location.
627 try:
628 if move_original and instance.source_path.is_file():
629 logger.info("Restoring previous original path")
630 shutil.move(instance.source_path, old_source_path)
632 if move_archive and instance.archive_path.is_file():
633 logger.info("Restoring previous archive path")
634 shutil.move(instance.archive_path, old_archive_path)
636 except Exception:
637 # This is fine, since:
638 # A: if we managed to move source from A to B, we will also
639 # manage to move it from B to A. If not, we have a serious
640 # issue that's going to get caught by the santiy checker.
641 # All files remain in place and will never be overwritten,
642 # so this is not the end of the world.
643 # B: if moving the original file failed, nothing has changed
644 # anyway.
645 pass
647 # restore old values on the instance
648 instance.filename = old_filename
649 instance.archive_filename = old_archive_filename
651 # finally, remove any empty sub folders. This will do nothing if
652 # something has failed above.
653 if not old_source_path.is_file(): 653 ↛ 659line 653 didn't jump to line 659 because the condition on line 653 was always true
654 delete_empty_directories(
655 Path(old_source_path).parent,
656 root=settings.ORIGINALS_DIR,
657 )
659 if instance.has_archive_version and not old_archive_path.is_file(): 659 ↛ 660line 659 didn't jump to line 660 because the condition on line 659 was never true
660 delete_empty_directories(
661 Path(old_archive_path).parent,
662 root=settings.ARCHIVE_DIR,
663 )
665 # Keep version files in sync with root
666 if instance.root_document_id is None: 666 ↛ exitline 666 didn't return from function 'update_filename_and_move_files' because the condition on line 666 was always true
667 for version_doc in Document.objects.filter(root_document_id=instance.pk).only( 667 ↛ 670line 667 didn't jump to line 670 because the loop on line 667 never started
668 "pk",
669 ):
670 update_filename_and_move_files(
671 Document,
672 version_doc,
673 )
676@shared_task
677def process_cf_select_update(custom_field: CustomField) -> None:
678 """
679 Update documents tied to a select custom field:
681 1. 'Select' custom field instances get their end-user value (e.g. in file names) from the select_options in extra_data,
682 which is contained in the custom field itself. So when the field is changed, we (may) need to update the file names
683 of all documents that have this custom field.
684 2. If a 'Select' field option was removed, we need to nullify the custom field instances that have the option.
685 """
686 select_options = {
687 option["id"]: option["label"]
688 for option in custom_field.extra_data.get("select_options", [])
689 }
691 # Clear select values that no longer exist
692 custom_field.fields.exclude(
693 value_select__in=select_options.keys(),
694 ).update(value_select=None)
696 for cf_instance in custom_field.fields.select_related("document").iterator():
697 # Update the filename and move files if necessary
698 update_filename_and_move_files(CustomFieldInstance, cf_instance)
701# should be disabled in /src/documents/management/commands/document_importer.py handle
702@receiver(models.signals.post_save, sender=CustomField)
703def check_paths_and_prune_custom_fields(
704 sender,
705 instance: CustomField,
706 **kwargs,
707) -> None:
708 """
709 When a custom field is updated, check if we need to update any documents. Done async to avoid slowing down the save operation.
710 """
711 if ( 711 ↛ 716line 711 didn't jump to line 716 because the condition on line 711 was never true
712 instance.data_type == CustomField.FieldDataType.SELECT
713 and instance.fields.count() > 0
714 and instance.extra_data
715 ): # Only select fields, for now
716 process_cf_select_update.apply_async(kwargs={"custom_field": instance})
719@receiver(models.signals.post_delete, sender=CustomField)
720def cleanup_custom_field_deletion(sender, instance: CustomField, **kwargs) -> None:
721 """
722 When a custom field is deleted, ensure no saved views reference it.
723 """
724 field_identifier = SavedView.DisplayFields.CUSTOM_FIELD % instance.pk
725 # remove field from display_fields of all saved views
726 for view in SavedView.objects.filter(display_fields__isnull=False).distinct(): 726 ↛ 727line 726 didn't jump to line 727 because the loop on line 726 never started
727 if field_identifier in view.display_fields:
728 logger.debug(
729 f"Removing custom field {instance} from view {view}",
730 )
731 view.display_fields.remove(field_identifier)
732 view.save()
734 # remove from sort_field of all saved views
735 views_with_sort_updated = SavedView.objects.filter(
736 sort_field=field_identifier,
737 ).update(
738 sort_field=SavedView.DisplayFields.CREATED,
739 )
740 if views_with_sort_updated > 0: 740 ↛ 741line 740 didn't jump to line 741 because the condition on line 740 was never true
741 logger.debug(
742 f"Removing custom field {instance} from sort field of {views_with_sort_updated} views",
743 )
746@receiver(models.signals.post_save, sender=Document)
747def update_llm_suggestions_cache(sender, instance, **kwargs):
748 """
749 Invalidate suggestions caches when a document is saved.
750 """
751 invalidate_suggestions_cache(instance.pk)
752 invalidate_llm_suggestions_cache(instance.pk)
755@receiver(models.signals.post_delete, sender=User)
756@receiver(models.signals.post_delete, sender=Group)
757def cleanup_user_deletion(sender, instance: User | Group, **kwargs) -> None:
758 """
759 When a user or group is deleted, remove non-cascading references.
760 At the moment, just the default permission settings in UiSettings.
761 """
762 # Remove the user permission settings e.g.
763 # DEFAULT_PERMS_OWNER: 'general-settings:permissions:default-owner',
764 # DEFAULT_PERMS_VIEW_USERS: 'general-settings:permissions:default-view-users',
765 # DEFAULT_PERMS_VIEW_GROUPS: 'general-settings:permissions:default-view-groups',
766 # DEFAULT_PERMS_EDIT_USERS: 'general-settings:permissions:default-edit-users',
767 # DEFAULT_PERMS_EDIT_GROUPS: 'general-settings:permissions:default-edit-groups',
768 for ui_settings in UiSettings.objects.all():
769 try:
770 permissions = ui_settings.settings.get("permissions", {})
771 updated = False
772 if isinstance(instance, User):
773 if permissions.get("default_owner") == instance.pk: 773 ↛ 774line 773 didn't jump to line 774 because the condition on line 773 was never true
774 permissions["default_owner"] = None
775 updated = True
776 if instance.pk in permissions.get("default_view_users", []): 776 ↛ 777line 776 didn't jump to line 777 because the condition on line 776 was never true
777 permissions["default_view_users"].remove(instance.pk)
778 updated = True
779 if instance.pk in permissions.get("default_change_users", []): 779 ↛ 780line 779 didn't jump to line 780 because the condition on line 779 was never true
780 permissions["default_change_users"].remove(instance.pk)
781 updated = True
782 elif isinstance(instance, Group): 782 ↛ 789line 782 didn't jump to line 789 because the condition on line 782 was always true
783 if instance.pk in permissions.get("default_view_groups", []): 783 ↛ 784line 783 didn't jump to line 784 because the condition on line 783 was never true
784 permissions["default_view_groups"].remove(instance.pk)
785 updated = True
786 if instance.pk in permissions.get("default_change_groups", []): 786 ↛ 787line 786 didn't jump to line 787 because the condition on line 786 was never true
787 permissions["default_change_groups"].remove(instance.pk)
788 updated = True
789 if updated: 789 ↛ 790line 789 didn't jump to line 790 because the condition on line 789 was never true
790 ui_settings.settings["permissions"] = permissions
791 ui_settings.save(update_fields=["settings"])
792 except Exception as e:
793 logger.error(
794 f"Error while cleaning up user {instance.pk} ({instance.username}) from ui_settings: {e}"
795 if isinstance(instance, User)
796 else f"Error while cleaning up group {instance.pk} ({instance.name}) from ui_settings: {e}",
797 )
800def add_to_index(sender, document, **kwargs) -> None:
801 from documents.search import get_backend
803 # A newly consumed version is not searchable on its own, its content
804 # becomes the effective_content of the root document
805 if document.root_document_id:
806 document = document.root_document
808 get_backend().add_or_update(document)
811def run_workflows_added(
812 sender,
813 document: Document,
814 logging_group: uuid.UUID | None = None,
815 original_file=None,
816 **kwargs,
817) -> None:
818 run_workflows(
819 trigger_type=WorkflowTrigger.WorkflowTriggerType.DOCUMENT_ADDED,
820 document=document,
821 logging_group=logging_group,
822 overrides=None,
823 original_file=original_file,
824 )
827def run_workflows_updated(
828 sender,
829 document: Document,
830 logging_group: uuid.UUID | None = None,
831 **kwargs,
832) -> None:
833 run_workflows(
834 trigger_type=WorkflowTrigger.WorkflowTriggerType.DOCUMENT_UPDATED,
835 document=document,
836 logging_group=logging_group,
837 )
840def send_websocket_document_updated(
841 sender,
842 document: Document,
843 **kwargs,
844) -> None:
845 # At this point, workflows may already have applied additional changes.
846 document.refresh_from_db()
848 from documents.data_models import DocumentMetadataOverrides
850 doc_overrides = DocumentMetadataOverrides.from_document(document)
852 with DocumentsStatusManager() as status_mgr:
853 status_mgr.send_document_updated(
854 document_id=document.id,
855 modified=DRF_DATETIME_FIELD.to_representation(document.modified),
856 owner_id=doc_overrides.owner_id,
857 users_can_view=doc_overrides.view_users,
858 groups_can_view=doc_overrides.view_groups,
859 )
862def run_workflows(
863 trigger_type: WorkflowTrigger.WorkflowTriggerType,
864 document: Document | ConsumableDocument,
865 workflow_to_run: Workflow | None = None,
866 logging_group: uuid.UUID | None = None,
867 overrides: DocumentMetadataOverrides | None = None,
868 original_file: Path | None = None,
869) -> tuple[DocumentMetadataOverrides, str] | None:
870 """
871 Execute workflows matching a document for the given trigger. When `overrides` is provided
872 (consumption flow), actions mutate that object and the function returns `(overrides, messages)`.
873 Otherwise actions mutate the actual document and return nothing.
875 Attachments for email/webhook actions use `original_file` when given, otherwise fall back to
876 `document.source_path` (Document) or `document.original_file` (ConsumableDocument).
878 Passing `workflow_to_run` skips the workflow query (currently only used by scheduled runs).
879 """
881 use_overrides = overrides is not None
883 if isinstance(document, Document) and document.root_document_id is not None:
884 logger.debug(
885 "Skipping workflow execution for version document %s",
886 document.pk,
887 )
888 return None
890 # Track whether the caller supplied original_file. When set explicitly (e.g. by
891 # run_workflows_added during consumption), it points at the staged file that has
892 # not yet been moved into its final storage location. This matters for password
893 # removal, which must read from the staged path rather than document.source_path.
894 caller_supplied_original_file = original_file is not None
895 if original_file is None:
896 original_file = (
897 document.source_path if not use_overrides else document.original_file
898 )
899 messages = []
901 workflows = get_workflows_for_trigger(trigger_type, workflow_to_run)
903 for workflow in workflows:
904 if not use_overrides:
905 if TYPE_CHECKING:
906 assert isinstance(document, Document)
907 try:
908 # This can be called from bulk_update_documents, which may be running multiple times
909 # Refresh this so the matching data is fresh and instance fields are re-freshed
910 # Otherwise, this instance might be behind and overwrite the work another process did
911 document.refresh_from_db()
912 except Document.DoesNotExist:
913 # Document was hard deleted by a previous workflow or another process
914 logger.info(
915 "Document no longer exists, skipping remaining workflows",
916 extra={"group": logging_group},
917 )
918 break
920 # Check if document was soft deleted (moved to trash)
921 if document.is_deleted:
922 logger.info(
923 "Document was moved to trash, skipping remaining workflows",
924 extra={"group": logging_group},
925 )
926 break
928 if matching.document_matches_workflow(document, workflow, trigger_type):
929 action: WorkflowAction
930 has_move_to_trash_action = False
931 for action in workflow.actions.order_by("order", "pk"):
932 message = f"Applying {action} from {workflow}"
933 if not use_overrides:
934 logger.info(message, extra={"group": logging_group})
935 else:
936 messages.append(message)
938 if action.type == WorkflowAction.WorkflowActionType.ASSIGNMENT:
939 if use_overrides and overrides:
940 apply_assignment_to_overrides(action, overrides)
941 else:
942 apply_assignment_to_document(
943 action,
944 document,
945 logging_group,
946 )
947 elif action.type == WorkflowAction.WorkflowActionType.REMOVAL:
948 if use_overrides and overrides:
949 apply_removal_to_overrides(action, overrides)
950 else:
951 apply_removal_to_document(action, document)
952 elif action.type == WorkflowAction.WorkflowActionType.EMAIL:
953 context = build_workflow_action_context(document, overrides)
954 execute_email_action(
955 action,
956 document,
957 context,
958 logging_group,
959 original_file,
960 trigger_type,
961 )
962 elif action.type == WorkflowAction.WorkflowActionType.WEBHOOK:
963 context = build_workflow_action_context(document, overrides)
964 execute_webhook_action(
965 action,
966 document,
967 context,
968 logging_group,
969 original_file,
970 )
971 elif action.type == WorkflowAction.WorkflowActionType.PASSWORD_REMOVAL:
972 execute_password_removal_action(
973 action,
974 document,
975 logging_group,
976 source_file=(
977 original_file if caller_supplied_original_file else None
978 ),
979 )
980 elif action.type == WorkflowAction.WorkflowActionType.MOVE_TO_TRASH:
981 has_move_to_trash_action = True
982 elif action.type == WorkflowAction.WorkflowActionType.REMOTE_OCR:
983 if use_overrides and overrides:
984 overrides.remote_ocr = True
985 else:
986 # If a workflow has a consumption trigger *and* another type,
987 # the document has already been parsed by the time the other one fires
988 logger.debug(
989 "Remote OCR action only applies to consumption "
990 "triggers, ignoring",
991 extra={"group": logging_group},
992 )
993 elif (
994 action.type
995 == WorkflowAction.WorkflowActionType.APPLY_AI_SUGGESTIONS
996 ):
997 if use_overrides:
998 # The document has not been parsed yet, so there is no
999 # content for the LLM to make suggestions from
1000 logger.debug(
1001 "Apply AI suggestions action does not apply to "
1002 "consumption triggers, ignoring",
1003 extra={"group": logging_group},
1004 )
1005 else:
1006 # Queued rather than run sync
1007 from documents.tasks import apply_ai_suggestions
1009 # kwargs so the PaperlessTask record can note the
1010 # document, see _extract_input_data
1011 apply_ai_suggestions.delay_on_commit(
1012 action_id=action.pk,
1013 document_id=document.pk,
1014 )
1016 if not use_overrides:
1017 # limit title to 128 characters
1018 document.title = document.title[:128]
1019 # Save only the fields that workflow actions can set directly.
1020 # Deliberately excludes filename and archive_filename — those are
1021 # managed exclusively by update_filename_and_move_files via the
1022 # post_save signal. Writing stale in-memory values here would revert
1023 # a concurrent update_filename_and_move_files DB write, leaving the
1024 # DB pointing at the old path while the file is already at the new
1025 # one (see: https://github.com/paperless-ngx/paperless-ngx/issues/12386).
1026 # modified has auto_now=True but is not auto-added when update_fields
1027 # is specified, so it must be listed explicitly.
1028 document.save(
1029 update_fields=[
1030 "title",
1031 "correspondent",
1032 "document_type",
1033 "storage_path",
1034 "owner",
1035 "modified",
1036 ],
1037 )
1039 WorkflowRun.objects.create(
1040 workflow=workflow,
1041 type=trigger_type,
1042 document=document if not use_overrides else None,
1043 )
1045 if has_move_to_trash_action:
1046 execute_move_to_trash_action(action, document, logging_group)
1048 if use_overrides:
1049 if TYPE_CHECKING:
1050 assert overrides is not None
1051 return overrides, "\n".join(messages)
1054# ---------------------------------------------------------------------------
1055# Task tracking -- Celery signal handlers
1056# ---------------------------------------------------------------------------
1058TRACKED_TASKS: dict[str, PaperlessTask.TaskType] = {
1059 "documents.tasks.consume_file": PaperlessTask.TaskType.CONSUME_FILE,
1060 "documents.tasks.train_classifier": PaperlessTask.TaskType.TRAIN_CLASSIFIER,
1061 "documents.tasks.sanity_check": PaperlessTask.TaskType.SANITY_CHECK,
1062 "documents.tasks.llmindex_index": PaperlessTask.TaskType.LLM_INDEX,
1063 "documents.tasks.empty_trash": PaperlessTask.TaskType.EMPTY_TRASH,
1064 "documents.tasks.check_scheduled_workflows": PaperlessTask.TaskType.CHECK_WORKFLOWS,
1065 "paperless_mail.tasks.process_mail_accounts": PaperlessTask.TaskType.MAIL_FETCH,
1066 "documents.tasks.bulk_update_documents": PaperlessTask.TaskType.BULK_UPDATE,
1067 "documents.tasks.update_document_content_maybe_archive_file": PaperlessTask.TaskType.REPROCESS_DOCUMENT,
1068 "documents.tasks.build_share_link_bundle": PaperlessTask.TaskType.BUILD_SHARE_LINK,
1069 "documents.bulk_edit.delete": PaperlessTask.TaskType.BULK_DELETE,
1070 "documents.tasks.apply_ai_suggestions": PaperlessTask.TaskType.APPLY_AI_SUGGESTIONS,
1071}
1073_CELERY_STATE_TO_STATUS: dict[str, PaperlessTask.Status] = {
1074 "SUCCESS": PaperlessTask.Status.SUCCESS,
1075 "FAILURE": PaperlessTask.Status.FAILURE,
1076 "REVOKED": PaperlessTask.Status.REVOKED,
1077}
1080def _extract_input_data(
1081 task_type: PaperlessTask.TaskType,
1082 task_kwargs: dict,
1083) -> dict:
1084 """Build the input_data dict stored on the PaperlessTask record.
1086 For consume_file tasks this includes the filename, MIME type, and any
1087 non-null overrides from the DocumentMetadataOverrides object. For
1088 mail_fetch tasks it captures the account_ids list. All other task
1089 types store no input data and return {}.
1090 """
1091 if task_type == PaperlessTask.TaskType.CONSUME_FILE:
1092 input_doc = task_kwargs.get("input_doc")
1093 overrides = task_kwargs.get("overrides")
1094 if input_doc is None: 1094 ↛ 1095line 1094 didn't jump to line 1095 because the condition on line 1094 was never true
1095 return {}
1096 data: dict = {
1097 "filename": input_doc.original_file.name,
1098 "mime_type": input_doc.mime_type,
1099 }
1100 if input_doc.original_path: # pragma: no cover 1100 ↛ 1101line 1100 didn't jump to line 1101 because the condition on line 1100 was never true
1101 data["source_path"] = str(input_doc.original_path)
1102 if input_doc.mailrule_id: # pragma: no cover 1102 ↛ 1103line 1102 didn't jump to line 1103 because the condition on line 1102 was never true
1103 data["mailrule_id"] = input_doc.mailrule_id
1104 if overrides: 1104 ↛ 1116line 1104 didn't jump to line 1116 because the condition on line 1104 was always true
1105 override_dict = {}
1106 for k, v in vars(overrides).items():
1107 if v is None or k.startswith("_"):
1108 continue
1109 if isinstance(v, datetime.date):
1110 v = v.isoformat()
1111 elif isinstance(v, Path): 1111 ↛ 1112line 1111 didn't jump to line 1112 because the condition on line 1111 was never true
1112 v = str(v)
1113 override_dict[k] = v
1114 if override_dict: 1114 ↛ 1116line 1114 didn't jump to line 1116 because the condition on line 1114 was always true
1115 data["overrides"] = override_dict
1116 return data
1118 if task_type == PaperlessTask.TaskType.MAIL_FETCH:
1119 account_ids = task_kwargs.get("account_ids")
1120 if account_ids is not None: 1120 ↛ 1122line 1120 didn't jump to line 1122 because the condition on line 1120 was always true
1121 return {"account_ids": account_ids}
1122 return {}
1124 if task_type == PaperlessTask.TaskType.APPLY_AI_SUGGESTIONS: 1124 ↛ 1125line 1124 didn't jump to line 1125 because the condition on line 1124 was never true
1125 document_id = task_kwargs.get("document_id")
1126 if document_id is not None:
1127 return {"document_id": document_id}
1128 return {}
1130 return {}
1133def _determine_trigger_source(
1134 headers: dict,
1135) -> PaperlessTask.TriggerSource:
1136 """Resolve the TriggerSource for a task being published to the broker.
1138 Reads the trigger_source header set by the caller; falls back to MANUAL
1139 when the header is absent or contains an unrecognised value.
1140 """
1141 header_source = headers.get("trigger_source")
1142 if header_source is not None: 1142 ↛ 1147line 1142 didn't jump to line 1147 because the condition on line 1142 was always true
1143 try:
1144 return PaperlessTask.TriggerSource(header_source)
1145 except ValueError:
1146 pass
1147 return PaperlessTask.TriggerSource.MANUAL
1150def _extract_owner_id(
1151 task_type: PaperlessTask.TaskType,
1152 task_kwargs: dict,
1153) -> int | None:
1154 """Return the owner_id from consume_file overrides, or None for all other task types."""
1155 if task_type != PaperlessTask.TaskType.CONSUME_FILE:
1156 return None
1157 overrides = task_kwargs.get("overrides")
1158 if overrides and hasattr(overrides, "owner_id"): 1158 ↛ 1160line 1158 didn't jump to line 1160 because the condition on line 1158 was always true
1159 return overrides.owner_id
1160 return None # pragma: no cover
1163@before_task_publish.connect
1164def before_task_publish_handler(
1165 sender=None,
1166 headers=None,
1167 body=None,
1168 **kwargs,
1169) -> None:
1170 """
1171 Creates the PaperlessTask record when the task is published to broker.
1173 https://docs.celeryq.dev/en/stable/userguide/signals.html#before-task-publish
1174 https://docs.celeryq.dev/en/stable/internals/protocol.html#version-2
1175 """
1176 if headers is None or body is None: 1176 ↛ 1177line 1176 didn't jump to line 1177 because the condition on line 1176 was never true
1177 return
1179 task_name = headers.get("task", "")
1180 task_type = TRACKED_TASKS.get(task_name)
1181 if task_type is None: 1181 ↛ 1182line 1181 didn't jump to line 1182 because the condition on line 1181 was never true
1182 return
1184 try:
1185 # Close stale connections without disrupting a transaction publishing a task
1186 for connection in connections.all(initialized_only=True):
1187 if not connection.in_atomic_block: 1187 ↛ 1186line 1187 didn't jump to line 1186 because the condition on line 1187 was always true
1188 connection.close_if_unusable_or_obsolete()
1190 _, task_kwargs, _ = body
1191 task_id = headers["id"]
1193 input_data = _extract_input_data(task_type, task_kwargs)
1194 trigger_source = _determine_trigger_source(headers)
1195 owner_id = _extract_owner_id(task_type, task_kwargs)
1197 # A retried task is republished with the same task_id, so this fires
1198 # again for it; get_or_create keeps the original PENDING record
1199 # instead of raising a duplicate-key IntegrityError on the retry.
1200 PaperlessTask.objects.get_or_create(
1201 task_id=task_id,
1202 defaults={
1203 "task_type": task_type,
1204 "trigger_source": trigger_source,
1205 "status": PaperlessTask.Status.PENDING,
1206 "input_data": input_data,
1207 "owner_id": owner_id,
1208 },
1209 )
1210 except Exception: # pragma: no cover
1211 logger.exception("Creating PaperlessTask failed")
1214@task_prerun.connect
1215def task_prerun_handler(sender=None, task_id=None, task=None, **kwargs) -> None:
1216 """
1217 Marks the task STARTED when execution begins on a worker.
1219 https://docs.celeryq.dev/en/stable/userguide/signals.html#task-prerun
1220 """
1221 if task_id is None: # pragma: no cover
1222 return
1223 if task and task.name not in TRACKED_TASKS:
1224 return
1225 try:
1226 close_old_connections()
1227 PaperlessTask.objects.filter(task_id=task_id).update(
1228 status=PaperlessTask.Status.STARTED,
1229 date_started=timezone.now(),
1230 )
1231 except Exception: # pragma: no cover
1232 logger.exception("Setting PaperlessTask started failed")
1235@task_postrun.connect
1236def task_postrun_handler(
1237 sender=None,
1238 task_id=None,
1239 task=None,
1240 retval=None,
1241 state=None,
1242 **kwargs,
1243) -> None:
1244 """
1245 Records task completion and result data for non-failure outcomes.
1247 Skips FAILURE states entirely, since task_failure_handler fires first
1248 and fully owns the failure path (status, date_done, duration, result_data).
1250 https://docs.celeryq.dev/en/stable/userguide/signals.html#task-postrun
1251 """
1252 if task_id is None: # pragma: no cover
1253 return
1254 if task and task.name not in TRACKED_TASKS:
1255 return
1256 try:
1257 close_old_connections()
1259 new_status = _CELERY_STATE_TO_STATUS.get(state, PaperlessTask.Status.FAILURE)
1260 if new_status == PaperlessTask.Status.FAILURE:
1261 return
1263 now = timezone.now()
1264 try:
1265 task_instance = PaperlessTask.objects.get(task_id=task_id)
1266 except PaperlessTask.DoesNotExist:
1267 return
1269 task_instance.status = new_status
1270 task_instance.date_done = now
1271 changed_fields = ["status", "date_done"]
1273 if task_instance.date_started:
1274 task_instance.duration_seconds = (
1275 now - task_instance.date_started
1276 ).total_seconds()
1277 changed_fields.append("duration_seconds")
1278 if task_instance.date_started and task_instance.date_created:
1279 task_instance.wait_time_seconds = (
1280 task_instance.date_started - task_instance.date_created
1281 ).total_seconds()
1282 changed_fields.append("wait_time_seconds")
1284 if isinstance(retval, dict):
1285 task_instance.result_data = retval
1286 changed_fields.append("result_data")
1287 if "duplicate_of" in retval:
1288 task_instance.status = PaperlessTask.Status.FAILURE
1289 changed_fields.append("status")
1291 task_instance.save(update_fields=changed_fields)
1292 except Exception: # pragma: no cover
1293 logger.exception("Updating PaperlessTask failed")
1296@task_failure.connect
1297def task_failure_handler(
1298 sender=None,
1299 task_id=None,
1300 exception=None,
1301 args=None,
1302 traceback=None,
1303 **kwargs,
1304) -> None:
1305 """
1306 Records failure details when a task raises an exception.
1308 Fully owns the FAILURE path. task_postrun_handler skips FAILURE
1309 states so there is no overlap.
1311 https://docs.celeryq.dev/en/stable/userguide/signals.html#task-failure
1312 """
1313 if task_id is None: # pragma: no cover
1314 return
1315 if sender and sender.name not in TRACKED_TASKS: # pragma: no cover
1316 return
1317 try:
1318 close_old_connections()
1320 result_data: dict = {
1321 "error_type": type(exception).__name__ if exception else "Unknown",
1322 "error_message": str(exception) if exception else "Unknown error",
1323 }
1324 if traceback:
1325 # billiard/celery pass a pre-formatted string instead of a real
1326 # traceback object when the worker process itself died (e.g.
1327 # WorkerLostError from a SIGILL) since there's no live traceback
1328 # to walk in that case.
1329 tb_str = (
1330 traceback
1331 if isinstance(traceback, str)
1332 else "".join(
1333 _tb.format_tb(traceback),
1334 )
1335 )
1336 result_data["traceback"] = tb_str[:5000]
1338 now = timezone.now()
1339 update_fields: dict = {
1340 "status": PaperlessTask.Status.FAILURE,
1341 "result_data": result_data,
1342 "date_done": now,
1343 }
1345 task_qs = PaperlessTask.objects.filter(task_id=task_id)
1346 task_instance = task_qs.values("date_started", "date_created").first()
1347 if task_instance:
1348 date_started = task_instance["date_started"]
1349 if date_started:
1350 update_fields["duration_seconds"] = (now - date_started).total_seconds()
1351 date_created = task_instance["date_created"]
1352 if date_started and date_created:
1353 update_fields["wait_time_seconds"] = (
1354 date_started - date_created
1355 ).total_seconds()
1356 task_qs.update(**update_fields)
1357 except Exception: # pragma: no cover
1358 logger.exception("Updating PaperlessTask on failure failed")
1361@task_revoked.connect
1362def task_revoked_handler(
1363 sender=None,
1364 request=None,
1365 *,
1366 terminated: bool = False,
1367 signum=None,
1368 expired: bool = False,
1369 **kwargs,
1370) -> None:
1371 """
1372 Marks the task REVOKED when it is cancelled before or during execution.
1374 This fires for tasks revoked while still queued (before task_prerun) as
1375 well as for tasks terminated mid-run. task_postrun does NOT fire for
1376 pre-start revocations, so this handler is the only way to move those
1377 records out of PENDING.
1379 https://docs.celeryq.dev/en/stable/userguide/signals.html#task-revoked
1380 """
1381 task_id = request.id if request else None
1382 if task_id is None: # pragma: no cover
1383 return
1384 if sender and sender.name not in TRACKED_TASKS: # pragma: no cover
1385 return
1386 try:
1387 close_old_connections()
1388 PaperlessTask.objects.filter(task_id=task_id).update(
1389 status=PaperlessTask.Status.REVOKED,
1390 date_done=timezone.now(),
1391 )
1392 except Exception: # pragma: no cover
1393 logger.exception("Updating PaperlessTask on revocation failed")
1396@worker_process_init.connect
1397def close_connection_pool_on_worker_init(**kwargs) -> None:
1398 """
1399 Close the DB connection pool for each Celery child process after it starts.
1401 This is necessary because the parent process parse the Django configuration,
1402 initializes connection pools then forks.
1404 Closing these pools after forking ensures child processes have a valid connection.
1405 """
1406 for conn in connections.all(initialized_only=True):
1407 if conn.alias == "default" and hasattr(conn, "pool") and conn.pool:
1408 conn.close_pool()
1411@worker_process_shutdown.connect
1412def close_connection_pool_on_worker_shutdown(**kwargs) -> None: # pragma: no cover
1413 """
1414 Close the DB connection pool when a Celery child process exits.
1416 With CELERY_WORKER_MAX_TASKS_PER_CHILD=1 each child is replaced after a
1417 single task. Without closing the pool on shutdown, its connections linger
1418 on the server until TCP keepalive reaps them, accumulating over time.
1419 """
1420 for conn in connections.all(initialized_only=True):
1421 if conn.alias == "default" and hasattr(conn, "pool") and conn.pool:
1422 conn.close_pool()
1425def add_or_update_document_in_llm_index(sender, document, **kwargs):
1426 """
1427 Add or update a document in the LLM index when it is created or updated.
1428 """
1429 if kwargs.get("skip_ai_index"):
1430 return
1431 ai_config = AIConfig()
1432 if ai_config.llm_index_enabled:
1433 from documents.tasks import update_document_in_llm_index
1435 update_document_in_llm_index.apply_async(kwargs={"document": document})
1438@receiver(models.signals.post_delete, sender=Document)
1439def delete_document_from_llm_index(
1440 sender: Any,
1441 instance: Document,
1442 **kwargs: Any,
1443) -> None:
1444 """
1445 Delete a document from the LLM index when it is deleted.
1446 """
1447 ai_config = AIConfig()
1448 if ai_config.llm_index_enabled: 1448 ↛ 1449line 1448 didn't jump to line 1449 because the condition on line 1448 was never true
1449 from documents.tasks import remove_document_from_llm_index
1451 remove_document_from_llm_index.apply_async(kwargs={"document": instance})