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

1from __future__ import annotations 

2 

3import datetime 

4import logging 

5import shutil 

6import traceback as _tb 

7from pathlib import Path 

8from typing import TYPE_CHECKING 

9from typing import Any 

10 

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 

31 

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 

71 

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

73 import uuid 

74 

75 from documents.classifier import DocumentClassifier 

76 from documents.data_models import ConsumableDocument 

77 from documents.data_models import DocumentMetadataOverrides 

78 

79logger = logging.getLogger("paperless.handlers") 

80DRF_DATETIME_FIELD = serializers.DateTimeField() 

81 

82 

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) 

94 

95 

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. 

109 

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

119 

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 

126 

127 potential_correspondents = matching.match_correspondents(document, classifier) 

128 potential_count = len(potential_correspondents) 

129 selected = potential_correspondents[0] if potential_correspondents else None 

130 

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 

145 

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

153 

154 return selected 

155 

156 

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. 

170 

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

180 

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 

187 

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 

191 

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 

206 

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

214 

215 return selected 

216 

217 

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. 

230 

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

233 

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

242 

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

261 

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

268 

269 current_tags = set(document.tags.all()) 

270 matched_tags = matching.match_tags(document, classifier) 

271 tags_to_add = set(matched_tags) - current_tags 

272 

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) 

279 

280 return tags_to_add, tags_to_remove 

281 

282 

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. 

296 

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

306 

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 

313 

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 

317 

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 

332 

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

340 

341 return selected 

342 

343 

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 

354 

355 while True: 

356 new_file_path = settings.EMPTY_TRASH_DIR / ( 

357 old_filebase + (f"_{counter:02}" if counter else "") + old_fileext 

358 ) 

359 

360 if new_file_path.exists(): 

361 counter += 1 

362 else: 

363 break 

364 

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 

378 

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

386 

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

399 

400 delete_empty_directories( 

401 Path(instance.source_path).parent, 

402 root=settings.ORIGINALS_DIR, 

403 ) 

404 

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 ) 

410 

411 

412class CannotMoveFilesException(Exception): 

413 pass 

414 

415 

416def _path_matches_checksum(path: Path, checksum: str | None) -> bool: 

417 if checksum is None or not path.is_file(): 

418 return False 

419 

420 return compute_checksum(path) == checksum 

421 

422 

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) 

429 

430 if not template: 

431 return False 

432 

433 return "custom_fields" in template 

434 

435 

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 

449 

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) 

458 

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) 

464 

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) 

470 

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 

481 

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

489 

490 old_filename = instance.filename 

491 old_source_path = instance.source_path 

492 move_original = False 

493 original_already_moved = False 

494 

495 old_archive_filename = instance.archive_filename 

496 old_archive_path = instance.archive_path 

497 move_archive = False 

498 archive_already_moved = False 

499 

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) 

509 

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 

530 

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 ) 

536 

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 

569 

570 instance.archive_filename = str(new_archive_filename) 

571 

572 move_archive = ( 

573 old_archive_filename != instance.archive_filename 

574 and not archive_already_moved 

575 ) 

576 else: 

577 move_archive = False 

578 

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 

585 

586 # Don't save() here to prevent infinite recursion. 

587 Document.objects.filter(pk=instance.pk).update(**updates) 

588 return 

589 

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) 

599 

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) 

609 

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) 

618 

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. 

625 

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) 

631 

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) 

635 

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 

646 

647 # restore old values on the instance 

648 instance.filename = old_filename 

649 instance.archive_filename = old_archive_filename 

650 

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 ) 

658 

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 ) 

664 

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 ) 

674 

675 

676@shared_task 

677def process_cf_select_update(custom_field: CustomField) -> None: 

678 """ 

679 Update documents tied to a select custom field: 

680 

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 } 

690 

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) 

695 

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) 

699 

700 

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

717 

718 

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

733 

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 ) 

744 

745 

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) 

753 

754 

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 ) 

798 

799 

800def add_to_index(sender, document, **kwargs) -> None: 

801 from documents.search import get_backend 

802 

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 

807 

808 get_backend().add_or_update(document) 

809 

810 

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 ) 

825 

826 

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 ) 

838 

839 

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

847 

848 from documents.data_models import DocumentMetadataOverrides 

849 

850 doc_overrides = DocumentMetadataOverrides.from_document(document) 

851 

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 ) 

860 

861 

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. 

874 

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

877 

878 Passing `workflow_to_run` skips the workflow query (currently only used by scheduled runs). 

879 """ 

880 

881 use_overrides = overrides is not None 

882 

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 

889 

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 = [] 

900 

901 workflows = get_workflows_for_trigger(trigger_type, workflow_to_run) 

902 

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 

919 

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 

927 

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) 

937 

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 

1008 

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 ) 

1015 

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 ) 

1038 

1039 WorkflowRun.objects.create( 

1040 workflow=workflow, 

1041 type=trigger_type, 

1042 document=document if not use_overrides else None, 

1043 ) 

1044 

1045 if has_move_to_trash_action: 

1046 execute_move_to_trash_action(action, document, logging_group) 

1047 

1048 if use_overrides: 

1049 if TYPE_CHECKING: 

1050 assert overrides is not None 

1051 return overrides, "\n".join(messages) 

1052 

1053 

1054# --------------------------------------------------------------------------- 

1055# Task tracking -- Celery signal handlers 

1056# --------------------------------------------------------------------------- 

1057 

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} 

1072 

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} 

1078 

1079 

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. 

1085 

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 

1117 

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 {} 

1123 

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 {} 

1129 

1130 return {} 

1131 

1132 

1133def _determine_trigger_source( 

1134 headers: dict, 

1135) -> PaperlessTask.TriggerSource: 

1136 """Resolve the TriggerSource for a task being published to the broker. 

1137 

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 

1148 

1149 

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 

1161 

1162 

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. 

1172 

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 

1178 

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 

1183 

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

1189 

1190 _, task_kwargs, _ = body 

1191 task_id = headers["id"] 

1192 

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) 

1196 

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

1212 

1213 

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. 

1218 

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

1233 

1234 

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. 

1246 

1247 Skips FAILURE states entirely, since task_failure_handler fires first 

1248 and fully owns the failure path (status, date_done, duration, result_data). 

1249 

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

1258 

1259 new_status = _CELERY_STATE_TO_STATUS.get(state, PaperlessTask.Status.FAILURE) 

1260 if new_status == PaperlessTask.Status.FAILURE: 

1261 return 

1262 

1263 now = timezone.now() 

1264 try: 

1265 task_instance = PaperlessTask.objects.get(task_id=task_id) 

1266 except PaperlessTask.DoesNotExist: 

1267 return 

1268 

1269 task_instance.status = new_status 

1270 task_instance.date_done = now 

1271 changed_fields = ["status", "date_done"] 

1272 

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

1283 

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

1290 

1291 task_instance.save(update_fields=changed_fields) 

1292 except Exception: # pragma: no cover 

1293 logger.exception("Updating PaperlessTask failed") 

1294 

1295 

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. 

1307 

1308 Fully owns the FAILURE path. task_postrun_handler skips FAILURE 

1309 states so there is no overlap. 

1310 

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

1319 

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] 

1337 

1338 now = timezone.now() 

1339 update_fields: dict = { 

1340 "status": PaperlessTask.Status.FAILURE, 

1341 "result_data": result_data, 

1342 "date_done": now, 

1343 } 

1344 

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

1359 

1360 

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. 

1373 

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. 

1378 

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

1394 

1395 

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. 

1400 

1401 This is necessary because the parent process parse the Django configuration, 

1402 initializes connection pools then forks. 

1403 

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

1409 

1410 

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. 

1415 

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

1423 

1424 

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 

1434 

1435 update_document_in_llm_index.apply_async(kwargs={"document": document}) 

1436 

1437 

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 

1450 

1451 remove_document_from_llm_index.apply_async(kwargs={"document": instance})