Coverage for documents/tasks.py: 25%

418 statements  

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

1import datetime 

2import logging 

3import shutil 

4import uuid 

5import zipfile 

6from collections.abc import Callable 

7from pathlib import Path 

8from tempfile import TemporaryDirectory 

9from tempfile import mkstemp 

10 

11from celery import Task 

12from celery import shared_task 

13from django.conf import settings 

14from django.contrib.contenttypes.models import ContentType 

15from django.db import models 

16from django.db import transaction 

17from django.db.models.signals import post_save 

18from django.utils import timezone 

19from filelock import FileLock 

20 

21from documents import sanity_checker 

22from documents.barcodes import BarcodePlugin 

23from documents.barcodes import read_barcode_values 

24from documents.bulk_download import ArchiveOnlyStrategy 

25from documents.bulk_download import OriginalsOnlyStrategy 

26from documents.caching import clear_document_caches 

27from documents.classifier import DocumentClassifier 

28from documents.classifier import load_classifier 

29from documents.consumer import AsnCheckPlugin 

30from documents.consumer import ConsumeFileDuplicateError 

31from documents.consumer import ConsumerPlugin 

32from documents.consumer import ConsumerPreflightPlugin 

33from documents.consumer import WorkflowTriggerPlugin 

34from documents.consumer import should_produce_archive 

35from documents.data_models import ConsumableDocument 

36from documents.data_models import ConsumeFileDuplicateResult 

37from documents.data_models import ConsumeFileStoppedResult 

38from documents.data_models import ConsumeFileSuccessResult 

39from documents.data_models import DocumentMetadataOverrides 

40from documents.data_models import StoredBarcode 

41from documents.double_sided import CollatePlugin 

42from documents.file_handling import create_source_path_directory 

43from documents.file_handling import generate_unique_filename 

44from documents.matching import prefilter_documents_by_workflowtrigger 

45from documents.models import Correspondent 

46from documents.models import CustomFieldInstance 

47from documents.models import Document 

48from documents.models import DocumentBarcode 

49from documents.models import DocumentType 

50from documents.models import PaperlessTask 

51from documents.models import ShareLink 

52from documents.models import ShareLinkBundle 

53from documents.models import StoragePath 

54from documents.models import Tag 

55from documents.models import WorkflowRun 

56from documents.models import WorkflowTrigger 

57from documents.plugins.base import ConsumeTaskPlugin 

58from documents.plugins.base import StopConsumeTaskError 

59from documents.plugins.helpers import ProgressManager 

60from documents.plugins.helpers import ProgressStatusOptions 

61from documents.sanity_checker import SanityCheckFailedException 

62from documents.search._backend import SearchIndexLockError 

63from documents.signals import document_updated 

64from documents.signals.handlers import cleanup_document_deletion 

65from documents.signals.handlers import run_workflows 

66from documents.signals.handlers import send_websocket_document_updated 

67from documents.utils import IterWrapper 

68from documents.utils import compute_checksum 

69from documents.utils import identity 

70from documents.versioning import annotate_effective_content 

71from documents.workflows.utils import get_workflows_for_trigger 

72from paperless.config import AIConfig 

73from paperless.config import BarcodeConfig 

74from paperless.config import RemoteOCRConfig 

75from paperless.logging import consume_task_id 

76from paperless.parsers import ParserContext 

77from paperless.parsers.registry import get_parser_registry 

78from paperless_ai.exceptions import LLMTimeoutError 

79from paperless_ai.indexing import llm_index_add_or_update_document 

80from paperless_ai.indexing import llm_index_remove_document 

81from paperless_ai.indexing import update_llm_index 

82 

83if settings.AUDIT_LOG_ENABLED: 83 ↛ 85line 83 didn't jump to line 85 because the condition on line 83 was always true

84 from auditlog.models import LogEntry 

85logger = logging.getLogger("paperless.tasks") 

86 

87 

88@shared_task 

89def index_optimize() -> None: 

90 logger.info( 

91 "index_optimize is a no-op — Tantivy manages segment merging automatically.", 

92 ) 

93 

94 

95@shared_task( 

96 bind=True, 

97 ignore_result=True, 

98 autoretry_for=(SearchIndexLockError,), 

99 max_retries=5, 

100 retry_backoff=60, 

101 retry_jitter=True, 

102) 

103def index_document(self, document_id: int) -> None: 

104 """ 

105 Deferred single-document index write. 

106 

107 Used as a self-healing fallback when add_or_update() exhausts its lock retry 

108 budget during high-concurrency consumption. Runs via batch_update() directly 

109 to avoid re-entering the deferred scheduling path in add_or_update(). 

110 

111 If the document was deleted before this task runs, it exits cleanly. 

112 """ 

113 from documents.search import get_backend 

114 

115 try: 

116 document = Document.objects.get(pk=document_id) 

117 except Document.DoesNotExist: 

118 logger.info( 

119 "index_document: document %d no longer exists; skipping", 

120 document_id, 

121 ) 

122 return 

123 with get_backend().batch_update() as batch: 

124 batch.add_or_update(document) 

125 

126 

127@shared_task( 

128 bind=True, 

129 ignore_result=True, 

130 autoretry_for=(SearchIndexLockError,), 

131 max_retries=5, 

132 retry_backoff=60, 

133 retry_jitter=True, 

134) 

135def remove_document_from_index(self, doc_id: int) -> None: 

136 """ 

137 Deferred single-document index removal. 

138 

139 Used as a self-healing fallback when remove() exhausts its lock retry budget. 

140 Operates only on the Tantivy index; no database lookup required. 

141 If the document has already been removed, the term-query delete is a no-op. 

142 """ 

143 from documents.search import get_backend 

144 

145 with get_backend().batch_update() as batch: 

146 batch.remove(doc_id) 

147 

148 

149@shared_task 

150def train_classifier( 

151 *, 

152 status_callback: Callable[[str], None] | None = None, 

153) -> str: 

154 if ( 

155 not Tag.objects.filter(matching_algorithm=Tag.MATCH_AUTO).exists() 

156 and not DocumentType.objects.filter(matching_algorithm=Tag.MATCH_AUTO).exists() 

157 and not Correspondent.objects.filter(matching_algorithm=Tag.MATCH_AUTO).exists() 

158 and not StoragePath.objects.filter(matching_algorithm=Tag.MATCH_AUTO).exists() 

159 ): 

160 result = "No automatic matching items, not training" 

161 logger.info(result) 

162 # Special case, items were once auto and trained, so remove the model 

163 # and prevent its use again 

164 if settings.MODEL_FILE.exists(): # pragma: no cover 

165 logger.info(f"Removing {settings.MODEL_FILE} so it won't be used") 

166 settings.MODEL_FILE.unlink() 

167 return result 

168 

169 classifier = load_classifier() 

170 

171 if not classifier: 

172 classifier = DocumentClassifier() 

173 

174 if classifier.train(status_callback=status_callback): 

175 logger.info( 

176 f"Saving updated classifier model to {settings.MODEL_FILE}...", 

177 ) 

178 classifier.save() 

179 return "Training completed successfully" 

180 else: 

181 logger.debug("Training data unchanged.") 

182 return "Training data unchanged" 

183 

184 

185@shared_task(bind=True) 

186def consume_file( 

187 self: Task, 

188 input_doc: ConsumableDocument, 

189 overrides: DocumentMetadataOverrides | None = None, 

190) -> ( 

191 ConsumeFileSuccessResult 

192 | ConsumeFileStoppedResult 

193 | ConsumeFileDuplicateResult 

194 | None 

195): 

196 token = consume_task_id.set((self.request.id or "")[:8]) 

197 try: 

198 # Default no overrides 

199 if overrides is None: 

200 overrides = DocumentMetadataOverrides() 

201 

202 plugins: list[type[ConsumeTaskPlugin]] = ( 

203 [ 

204 ConsumerPreflightPlugin, 

205 ConsumerPlugin, 

206 ] 

207 if input_doc.root_document_id is not None 

208 else [ 

209 ConsumerPreflightPlugin, 

210 AsnCheckPlugin, 

211 CollatePlugin, 

212 BarcodePlugin, 

213 AsnCheckPlugin, # Re-run ASN check after barcode reading 

214 WorkflowTriggerPlugin, 

215 ConsumerPlugin, 

216 ] 

217 ) 

218 

219 with ( 

220 ProgressManager( 

221 overrides.filename or input_doc.original_file.name, 

222 self.request.id, 

223 ) as status_mgr, 

224 TemporaryDirectory(dir=settings.SCRATCH_DIR) as tmp_dir, 

225 ): 

226 tmp_dir = Path(tmp_dir) 

227 msg = None 

228 for plugin_class in plugins: 

229 plugin_name = plugin_class.NAME 

230 

231 plugin = plugin_class( 

232 input_doc, 

233 overrides, 

234 status_mgr, 

235 tmp_dir, 

236 self.request.id, 

237 ) 

238 

239 if not plugin.able_to_run: 

240 logger.debug(f"Skipping plugin {plugin_name}") 

241 continue 

242 

243 try: 

244 logger.debug(f"Executing plugin {plugin_name}") 

245 plugin.setup() 

246 

247 msg = plugin.run() 

248 

249 if msg is not None: 

250 logger.info(f"{plugin_name} completed with: {msg}") 

251 else: 

252 logger.info(f"{plugin_name} completed with no message") 

253 

254 overrides = plugin.metadata 

255 

256 except StopConsumeTaskError as e: 

257 logger.info(f"{plugin_name} requested task exit: {e.message}") 

258 return ConsumeFileStoppedResult(reason=e.message) 

259 

260 except ConsumeFileDuplicateError as e: 

261 logger.info(f"{plugin_name} rejected duplicate: {e}") 

262 return ConsumeFileDuplicateResult( 

263 duplicate_of=e.duplicate_id, 

264 duplicate_in_trash=e.in_trash, 

265 ) 

266 

267 except Exception as e: 

268 logger.exception(f"{plugin_name} failed: {e}") 

269 status_mgr.send_progress( 

270 ProgressStatusOptions.FAILED, 

271 f"{e}", 

272 100, 

273 100, 

274 ) 

275 raise 

276 

277 finally: 

278 plugin.cleanup() 

279 

280 return msg 

281 finally: 

282 consume_task_id.reset(token) 

283 

284 

285@shared_task 

286def sanity_check(*, raise_on_error: bool = True) -> str: 

287 messages = sanity_checker.check_sanity() 

288 messages.log_messages() 

289 

290 if not messages.has_error and not messages.has_warning and not messages.has_info: 

291 return "No issues detected." 

292 

293 parts: list[str] = [] 

294 if messages.document_error_count: 

295 parts.append(f"{messages.document_error_count} document(s) with errors") 

296 if messages.document_warning_count: 

297 parts.append(f"{messages.document_warning_count} document(s) with warnings") 

298 if messages.document_info_count: 

299 parts.append(f"{messages.document_info_count} document(s) with infos") 

300 if messages.global_warning_count: 

301 parts.append(f"{messages.global_warning_count} global warning(s)") 

302 

303 summary = ", ".join(parts) + " found." 

304 

305 if messages.has_error: 

306 message = summary + " Check logs for details." 

307 if raise_on_error: 

308 raise SanityCheckFailedException(message) 

309 return message 

310 

311 return summary 

312 

313 

314@shared_task 

315def bulk_update_documents(document_ids) -> None: 

316 from documents.search import get_backend 

317 

318 document_ids = list(document_ids) 

319 # Annotated so the signal handlers below (e.g. matching) don't query the 

320 # versions of each document. Indexing re-queries and re-annotates its own 

321 # copy via add_or_update_ids() below, after these signals (and any 

322 # workflow they trigger) have had a chance to mutate the documents. 

323 documents = annotate_effective_content( 

324 Document.objects.filter(id__in=document_ids), 

325 ) 

326 

327 for doc in documents: 

328 clear_document_caches(doc.pk) 

329 document_updated.send( 

330 sender=None, 

331 document=doc, 

332 logging_group=uuid.uuid4(), 

333 skip_ai_index=True, # bulk path calls update_llm_index once below 

334 ) 

335 post_save.send(Document, instance=doc, created=False) 

336 

337 with get_backend().batch_update() as batch: 

338 batch.add_or_update_ids(document_ids) 

339 

340 ai_config = AIConfig() 

341 if ai_config.llm_index_enabled: 

342 update_llm_index( 

343 rebuild=False, 

344 document_ids=document_ids, 

345 ) 

346 

347 

348def _read_barcodes_for_reprocess(document: Document) -> list[StoredBarcode] | None: 

349 """ 

350 Reads the barcodes of the original again, e.g. for documents consumed 

351 before storing them was enabled. Returns None if they should be left as 

352 they are: storing is off, the file can't be scanned with the current 

353 settings, or the scan failed. 

354 """ 

355 barcode_settings = BarcodeConfig() 

356 if not barcode_settings.barcode_store_values: 

357 return None 

358 try: 

359 with TemporaryDirectory(dir=settings.SCRATCH_DIR) as tmpdir: 

360 return read_barcode_values( 

361 document.source_path, 

362 document.mime_type, 

363 barcode_settings, 

364 Path(tmpdir), 

365 ) 

366 except Exception as e: 

367 logger.warning(f"Could not read barcodes of document {document}: {e}") 

368 return None 

369 

370 

371@shared_task 

372def update_document_content_maybe_archive_file( 

373 document_id, 

374 *, 

375 remote_ocr: bool = False, 

376) -> None: 

377 """ 

378 Re-creates OCR content and thumbnail for a document, and archive file if 

379 it exists. 

380 

381 Remote OCR is used only when the engine is configured to handle everything 

382 or if explicitly asked for via ``remote_ocr``. 

383 """ 

384 document = Document.objects.get(id=document_id) 

385 

386 mime_type = document.mime_type 

387 

388 parser_class = get_parser_registry().get_parser_for_file( 

389 mime_type, 

390 document.original_filename or "", 

391 document.source_path, 

392 allow_remote=remote_ocr or RemoteOCRConfig().remote_ocr_by_default, 

393 ) 

394 

395 if not parser_class: 

396 logger.error( 

397 f"No parser found for mime type {mime_type}, cannot " 

398 f"archive document {document} (ID: {document_id})", 

399 ) 

400 return 

401 

402 with parser_class() as parser: 

403 parser.configure(ParserContext()) 

404 

405 try: 

406 produce_archive = should_produce_archive( 

407 parser, 

408 mime_type, 

409 document.source_path, 

410 ) 

411 parser.parse( 

412 document.source_path, 

413 mime_type, 

414 produce_archive=produce_archive, 

415 ) 

416 

417 barcodes = _read_barcodes_for_reprocess(document) 

418 

419 thumbnail = parser.get_thumbnail(document.source_path, mime_type) 

420 

421 with transaction.atomic(): 

422 oldDocument = Document.objects.get(pk=document.pk) 

423 if parser.get_archive_path(): 

424 checksum = compute_checksum(parser.get_archive_path()) 

425 # I'm going to save first so that in case the file move 

426 # fails, the database is rolled back. 

427 # We also don't use save() since that triggers the filehandling 

428 # logic, and we don't want that yet (file not yet in place) 

429 document.archive_filename = generate_unique_filename( 

430 document, 

431 archive_filename=True, 

432 ) 

433 Document.objects.filter(pk=document.pk).update( 

434 archive_checksum=checksum, 

435 content=parser.get_text(), 

436 archive_filename=document.archive_filename, 

437 ) 

438 newDocument = Document.objects.get(pk=document.pk) 

439 if settings.AUDIT_LOG_ENABLED: 

440 LogEntry.objects.log_create( 

441 instance=oldDocument, 

442 changes={ 

443 "content": [oldDocument.content, newDocument.content], 

444 "archive_checksum": [ 

445 oldDocument.archive_checksum, 

446 newDocument.archive_checksum, 

447 ], 

448 "archive_filename": [ 

449 oldDocument.archive_filename, 

450 newDocument.archive_filename, 

451 ], 

452 }, 

453 additional_data={ 

454 "reason": "Update document content", 

455 }, 

456 action=LogEntry.Action.UPDATE, 

457 ) 

458 else: 

459 Document.objects.filter(pk=document.pk).update( 

460 content=parser.get_text(), 

461 ) 

462 

463 if settings.AUDIT_LOG_ENABLED: 

464 LogEntry.objects.log_create( 

465 instance=oldDocument, 

466 changes={ 

467 "content": [oldDocument.content, parser.get_text()], 

468 }, 

469 additional_data={ 

470 "reason": "Update document content", 

471 }, 

472 action=LogEntry.Action.UPDATE, 

473 ) 

474 

475 if barcodes is not None: 

476 document.barcodes.all().delete() 

477 DocumentBarcode.objects.bulk_create( 

478 DocumentBarcode(document=document, **barcode) 

479 for barcode in barcodes 

480 ) 

481 # metadata_etag includes modified 

482 Document.objects.filter(pk=document.pk).update( 

483 modified=timezone.now(), 

484 ) 

485 

486 with FileLock(settings.MEDIA_LOCK): 

487 if parser.get_archive_path(): 

488 create_source_path_directory(document.archive_path) 

489 shutil.move(parser.get_archive_path(), document.archive_path) 

490 shutil.move(thumbnail, document.thumbnail_path) 

491 

492 document.refresh_from_db() 

493 logger.info( 

494 f"Updating index for document {document_id} ({document.archive_checksum})", 

495 ) 

496 from documents.search import get_backend 

497 

498 get_backend().add_or_update(document) 

499 

500 ai_config = AIConfig() 

501 if ai_config.llm_index_enabled: 

502 llm_index_add_or_update_document(document) 

503 

504 clear_document_caches(document.pk) 

505 

506 except Exception: 

507 logger.exception( 

508 f"Error while parsing document {document} (ID: {document_id})", 

509 ) 

510 

511 

512@shared_task 

513def empty_trash(doc_ids=None) -> None: 

514 if doc_ids is None: 514 ↛ 515line 514 didn't jump to line 515 because the condition on line 514 was never true

515 logger.info("Emptying trash of all expired documents") 

516 documents = ( 

517 Document.deleted_objects.filter(id__in=doc_ids) 

518 if doc_ids is not None 

519 else Document.deleted_objects.filter( 

520 deleted_at__lt=timezone.localtime(timezone.now()) 

521 - datetime.timedelta( 

522 days=settings.EMPTY_TRASH_DELAY, 

523 ), 

524 ) 

525 ) 

526 

527 try: 

528 deleted_document_ids = list(documents.values_list("id", flat=True)) 

529 # Temporarily connect the cleanup handler 

530 models.signals.post_delete.connect(cleanup_document_deletion, sender=Document) 

531 documents.delete() # this is effectively a hard delete 

532 logger.info(f"Deleted {len(deleted_document_ids)} documents from trash") 

533 

534 if settings.AUDIT_LOG_ENABLED: 534 ↛ 543line 534 didn't jump to line 543 because the condition on line 534 was always true

535 # Delete the audit log entries for documents that dont exist anymore 

536 LogEntry.objects.filter( 

537 content_type=ContentType.objects.get_for_model(Document), 

538 object_id__in=deleted_document_ids, 

539 ).delete() 

540 except Exception as e: # pragma: no cover 

541 logger.exception(f"Error while emptying trash: {e}") 

542 finally: 

543 models.signals.post_delete.disconnect( 

544 cleanup_document_deletion, 

545 sender=Document, 

546 ) 

547 

548 

549@shared_task 

550def check_scheduled_workflows() -> None: 

551 """ 

552 Check and run all enabled scheduled workflows. 

553 

554 Scheduled triggers are evaluated based on a target date field (e.g. added, created, modified, or a custom date field), 

555 combined with a day offset: 

556 - Positive offsets mean the workflow should trigger AFTER the specified date (e.g., offset = +7 → trigger 7 days after) 

557 - Negative offsets mean the workflow should trigger BEFORE the specified date (e.g., offset = -7 → trigger 7 days before) 

558 

559 Once a document satisfies this condition, and recurring/non-recurring constraints are met, the workflow is run. 

560 """ 

561 scheduled_workflows = get_workflows_for_trigger( 

562 WorkflowTrigger.WorkflowTriggerType.SCHEDULED, 

563 ) 

564 if scheduled_workflows.count() > 0: 

565 logger.debug(f"Checking {len(scheduled_workflows)} scheduled workflows") 

566 now = timezone.now() 

567 for workflow in scheduled_workflows: 

568 schedule_triggers = workflow.triggers.filter( 

569 type=WorkflowTrigger.WorkflowTriggerType.SCHEDULED, 

570 ) 

571 trigger: WorkflowTrigger 

572 for trigger in schedule_triggers: 

573 documents = Document.objects.none() 

574 offset_td = datetime.timedelta(days=trigger.schedule_offset_days) 

575 threshold = now - offset_td 

576 logger.debug( 

577 f"Trigger {trigger.id}: checking if (date + {offset_td}) <= now ({now})", 

578 ) 

579 

580 match trigger.schedule_date_field: 

581 case WorkflowTrigger.ScheduleDateField.ADDED: 

582 documents = Document.objects.filter( 

583 root_document__isnull=True, 

584 added__lte=threshold, 

585 ) 

586 

587 case WorkflowTrigger.ScheduleDateField.CREATED: 

588 documents = Document.objects.filter( 

589 root_document__isnull=True, 

590 created__lte=threshold, 

591 ) 

592 

593 case WorkflowTrigger.ScheduleDateField.MODIFIED: 

594 documents = Document.objects.filter( 

595 root_document__isnull=True, 

596 modified__lte=threshold, 

597 ) 

598 

599 case WorkflowTrigger.ScheduleDateField.CUSTOM_FIELD: 

600 # cap earliest date to avoid massive scans 

601 earliest_date = now - datetime.timedelta(days=365) 

602 if offset_td.days < -365: 

603 logger.warning( 

604 f"Trigger {trigger.id} has large negative offset ({offset_td.days}), " 

605 f"limiting earliest scan date to {earliest_date}", 

606 ) 

607 

608 cf_filter_kwargs = { 

609 "field": trigger.schedule_date_custom_field, 

610 "value_date__isnull": False, 

611 "value_date__lte": threshold, 

612 "value_date__gte": earliest_date, 

613 } 

614 

615 recent_cf_instances = CustomFieldInstance.objects.filter( 

616 **cf_filter_kwargs, 

617 ) 

618 

619 matched_ids = [ 

620 cfi.document_id 

621 for cfi in recent_cf_instances 

622 if cfi.value_date 

623 and ( 

624 timezone.make_aware( 

625 datetime.datetime.combine( 

626 cfi.value_date, 

627 datetime.time.min, 

628 ), 

629 ) 

630 + offset_td 

631 <= now 

632 ) 

633 ] 

634 

635 documents = Document.objects.filter( 

636 root_document__isnull=True, 

637 id__in=matched_ids, 

638 ) 

639 

640 if documents.exists(): 

641 documents = prefilter_documents_by_workflowtrigger( 

642 documents, 

643 trigger, 

644 ) 

645 

646 if documents.exists(): 

647 logger.debug( 

648 f"Found {documents.count()} documents for trigger {trigger}", 

649 ) 

650 for document in documents: 

651 workflow_runs = WorkflowRun.objects.filter( 

652 document=document, 

653 type=WorkflowTrigger.WorkflowTriggerType.SCHEDULED, 

654 workflow=workflow, 

655 ).order_by("-run_at") 

656 if not trigger.schedule_is_recurring and workflow_runs.exists(): 

657 logger.debug( 

658 f"Skipping document {document} for non-recurring workflow {workflow} as it has already been run", 

659 ) 

660 continue 

661 

662 if ( 

663 trigger.schedule_is_recurring 

664 and workflow_runs.exists() 

665 and ( 

666 workflow_runs.first().run_at 

667 > now 

668 - datetime.timedelta( 

669 days=trigger.schedule_recurring_interval_days, 

670 ) 

671 ) 

672 ): 

673 # schedule is recurring but the last run was within the number of recurring interval days 

674 logger.debug( 

675 f"Skipping document {document} for recurring workflow {workflow} as the last run was within the recurring interval", 

676 ) 

677 continue 

678 run_workflows( 

679 trigger_type=WorkflowTrigger.WorkflowTriggerType.SCHEDULED, 

680 workflow_to_run=workflow, 

681 document=document, 

682 ) 

683 # Scheduled workflows dont send document_updated signal, so send a websocket update here to ensure clients are updated 

684 send_websocket_document_updated( 

685 sender=None, 

686 document=document, 

687 ) 

688 

689 

690def update_document_parent_tags(tag: Tag, new_parent: Tag) -> None: 

691 """ 

692 When a tag's parent changes, ensure all documents containing the tag also have 

693 the parent tag (and its ancestors) applied. 

694 """ 

695 doc_tag_relationship = Document.tags.through 

696 

697 doc_ids: list[int] = list( 

698 Document.objects.filter(tags=tag).values_list("pk", flat=True), 

699 ) 

700 

701 if not doc_ids: 

702 return 

703 

704 parent_ids = [new_parent.id, *new_parent.get_ancestors_pks()] 

705 

706 parent_ids = list(dict.fromkeys(parent_ids)) 

707 

708 existing_pairs = set( 

709 doc_tag_relationship.objects.filter( 

710 document_id__in=doc_ids, 

711 tag_id__in=parent_ids, 

712 ).values_list("document_id", "tag_id"), 

713 ) 

714 

715 to_create: list = [] 

716 affected: set[int] = set() 

717 

718 for doc_id in doc_ids: 

719 for parent_id in parent_ids: 

720 if (doc_id, parent_id) in existing_pairs: 

721 continue 

722 

723 to_create.append( 

724 doc_tag_relationship(document_id=doc_id, tag_id=parent_id), 

725 ) 

726 affected.add(doc_id) 

727 

728 if to_create: 

729 doc_tag_relationship.objects.bulk_create( 

730 to_create, 

731 ignore_conflicts=True, 

732 ) 

733 

734 if affected: 

735 bulk_update_documents.apply_async( 

736 kwargs={"document_ids": list(affected)}, 

737 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM}, 

738 ) 

739 

740 

741@shared_task 

742def llmindex_index( 

743 *, 

744 iter_wrapper: IterWrapper[Document] = identity, 

745 rebuild: bool = False, 

746) -> str | None: 

747 ai_config = AIConfig() 

748 if not ai_config.llm_index_enabled: # pragma: no cover 

749 logger.info("LLM index is disabled, skipping update.") 

750 return None 

751 

752 from paperless_ai.indexing import update_llm_index 

753 

754 return update_llm_index( 

755 iter_wrapper=iter_wrapper, 

756 rebuild=rebuild, 

757 ) 

758 

759 

760@shared_task( 

761 bind=True, 

762 autoretry_for=(LLMTimeoutError,), 

763 max_retries=3, 

764 retry_backoff=60, 

765 retry_backoff_max=600, 

766 retry_jitter=True, 

767) 

768def apply_ai_suggestions(self, action_id: int, document_id: int) -> None: 

769 """ 

770 Deferred "apply AI suggestions" workflow action. 

771 """ 

772 from documents.models import WorkflowAction 

773 from documents.workflows.ai import apply_ai_suggestions_to_document 

774 

775 try: 

776 action = WorkflowAction.objects.get(pk=action_id) 

777 document = Document.objects.select_related("owner").get(pk=document_id) 

778 except (WorkflowAction.DoesNotExist, Document.DoesNotExist): 

779 logger.warning( 

780 "Workflow action %s or document %s no longer exists, " 

781 "not applying AI suggestions", 

782 action_id, 

783 document_id, 

784 ) 

785 return 

786 

787 if not apply_ai_suggestions_to_document(action, document): 

788 return 

789 

790 # No document_updated signal to avoid loop 

791 clear_document_caches(document.pk) 

792 index_document.delay(document.pk) 

793 

794 ai_config = AIConfig() 

795 if ai_config.llm_index_enabled: 

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

797 

798 

799@shared_task 

800def update_document_in_llm_index(document) -> None: 

801 llm_index_add_or_update_document(document) 

802 

803 

804@shared_task 

805def remove_document_from_llm_index(document) -> None: 

806 llm_index_remove_document(document) 

807 

808 

809@shared_task 

810def build_share_link_bundle(bundle_id: int) -> None: 

811 try: 

812 bundle = ( 

813 ShareLinkBundle.objects.filter(pk=bundle_id) 

814 .prefetch_related("documents") 

815 .get() 

816 ) 

817 except ShareLinkBundle.DoesNotExist: 

818 logger.warning("Share link bundle %s no longer exists.", bundle_id) 

819 return 

820 

821 bundle.remove_file() 

822 bundle.status = ShareLinkBundle.Status.PROCESSING 

823 bundle.last_error = None 

824 bundle.size_bytes = None 

825 bundle.built_at = None 

826 bundle.file_path = "" 

827 bundle.save( 

828 update_fields=[ 

829 "status", 

830 "last_error", 

831 "size_bytes", 

832 "built_at", 

833 "file_path", 

834 ], 

835 ) 

836 

837 documents = list(bundle.documents.all().order_by("pk")) 

838 

839 _, temp_zip_path_str = mkstemp(suffix=".zip", dir=settings.SCRATCH_DIR) 

840 temp_zip_path = Path(temp_zip_path_str) 

841 

842 try: 

843 strategy_class = ( 

844 ArchiveOnlyStrategy 

845 if bundle.file_version == ShareLink.FileVersion.ARCHIVE 

846 else OriginalsOnlyStrategy 

847 ) 

848 with zipfile.ZipFile(temp_zip_path, "w", zipfile.ZIP_DEFLATED) as zipf: 

849 strategy = strategy_class(zipf) 

850 for document in documents: 

851 strategy.add_document(document) 

852 

853 output_dir = settings.SHARE_LINK_BUNDLE_DIR 

854 output_dir.mkdir(parents=True, exist_ok=True) 

855 final_path = (output_dir / f"{bundle.slug}.zip").resolve() 

856 if final_path.exists(): 

857 final_path.unlink() 

858 shutil.move(temp_zip_path, final_path) 

859 

860 bundle.file_path = f"{bundle.slug}.zip" 

861 bundle.size_bytes = final_path.stat().st_size 

862 bundle.status = ShareLinkBundle.Status.READY 

863 bundle.built_at = timezone.now() 

864 bundle.last_error = None 

865 bundle.save( 

866 update_fields=[ 

867 "file_path", 

868 "size_bytes", 

869 "status", 

870 "built_at", 

871 "last_error", 

872 ], 

873 ) 

874 logger.info("Built share link bundle %s", bundle.pk) 

875 except Exception as exc: 

876 logger.exception( 

877 "Failed to build share link bundle %s: %s", 

878 bundle_id, 

879 exc, 

880 ) 

881 bundle.status = ShareLinkBundle.Status.FAILED 

882 bundle.last_error = { 

883 "bundle_id": bundle_id, 

884 "exception_type": exc.__class__.__name__, 

885 "message": str(exc), 

886 "timestamp": timezone.now().isoformat(), 

887 } 

888 bundle.save(update_fields=["status", "last_error"]) 

889 try: 

890 temp_zip_path.unlink() 

891 except OSError: 

892 pass 

893 raise 

894 finally: 

895 try: 

896 temp_zip_path.unlink(missing_ok=True) 

897 except OSError: 

898 pass 

899 

900 

901@shared_task 

902def cleanup_expired_share_link_bundles() -> None: 

903 now = timezone.now() 

904 expired_qs = ShareLinkBundle.objects.filter( 

905 expiration__isnull=False, 

906 expiration__lt=now, 

907 ) 

908 count = 0 

909 for bundle in expired_qs.iterator(): 

910 count += 1 

911 try: 

912 bundle.delete() 

913 except Exception as exc: 

914 logger.warning( 

915 "Failed to delete expired share link bundle %s: %s", 

916 bundle.pk, 

917 exc, 

918 ) 

919 if count: 

920 logger.info("Deleted %s expired share link bundle(s)", count)