Coverage for documents/tasks.py: 25%
418 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
1import datetime
2import logging
3import shutil
4import uuid
5import zipfile
6from collections.abc import Callable
7from pathlib import Path
8from tempfile import TemporaryDirectory
9from tempfile import mkstemp
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
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
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")
88@shared_task
89def index_optimize() -> None:
90 logger.info(
91 "index_optimize is a no-op — Tantivy manages segment merging automatically.",
92 )
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.
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().
111 If the document was deleted before this task runs, it exits cleanly.
112 """
113 from documents.search import get_backend
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)
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.
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
145 with get_backend().batch_update() as batch:
146 batch.remove(doc_id)
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
169 classifier = load_classifier()
171 if not classifier:
172 classifier = DocumentClassifier()
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"
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()
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 )
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
231 plugin = plugin_class(
232 input_doc,
233 overrides,
234 status_mgr,
235 tmp_dir,
236 self.request.id,
237 )
239 if not plugin.able_to_run:
240 logger.debug(f"Skipping plugin {plugin_name}")
241 continue
243 try:
244 logger.debug(f"Executing plugin {plugin_name}")
245 plugin.setup()
247 msg = plugin.run()
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")
254 overrides = plugin.metadata
256 except StopConsumeTaskError as e:
257 logger.info(f"{plugin_name} requested task exit: {e.message}")
258 return ConsumeFileStoppedResult(reason=e.message)
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 )
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
277 finally:
278 plugin.cleanup()
280 return msg
281 finally:
282 consume_task_id.reset(token)
285@shared_task
286def sanity_check(*, raise_on_error: bool = True) -> str:
287 messages = sanity_checker.check_sanity()
288 messages.log_messages()
290 if not messages.has_error and not messages.has_warning and not messages.has_info:
291 return "No issues detected."
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)")
303 summary = ", ".join(parts) + " found."
305 if messages.has_error:
306 message = summary + " Check logs for details."
307 if raise_on_error:
308 raise SanityCheckFailedException(message)
309 return message
311 return summary
314@shared_task
315def bulk_update_documents(document_ids) -> None:
316 from documents.search import get_backend
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 )
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)
337 with get_backend().batch_update() as batch:
338 batch.add_or_update_ids(document_ids)
340 ai_config = AIConfig()
341 if ai_config.llm_index_enabled:
342 update_llm_index(
343 rebuild=False,
344 document_ids=document_ids,
345 )
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
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.
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)
386 mime_type = document.mime_type
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 )
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
402 with parser_class() as parser:
403 parser.configure(ParserContext())
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 )
417 barcodes = _read_barcodes_for_reprocess(document)
419 thumbnail = parser.get_thumbnail(document.source_path, mime_type)
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 )
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 )
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 )
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)
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
498 get_backend().add_or_update(document)
500 ai_config = AIConfig()
501 if ai_config.llm_index_enabled:
502 llm_index_add_or_update_document(document)
504 clear_document_caches(document.pk)
506 except Exception:
507 logger.exception(
508 f"Error while parsing document {document} (ID: {document_id})",
509 )
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 )
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")
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 )
549@shared_task
550def check_scheduled_workflows() -> None:
551 """
552 Check and run all enabled scheduled workflows.
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)
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 )
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 )
587 case WorkflowTrigger.ScheduleDateField.CREATED:
588 documents = Document.objects.filter(
589 root_document__isnull=True,
590 created__lte=threshold,
591 )
593 case WorkflowTrigger.ScheduleDateField.MODIFIED:
594 documents = Document.objects.filter(
595 root_document__isnull=True,
596 modified__lte=threshold,
597 )
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 )
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 }
615 recent_cf_instances = CustomFieldInstance.objects.filter(
616 **cf_filter_kwargs,
617 )
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 ]
635 documents = Document.objects.filter(
636 root_document__isnull=True,
637 id__in=matched_ids,
638 )
640 if documents.exists():
641 documents = prefilter_documents_by_workflowtrigger(
642 documents,
643 trigger,
644 )
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
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 )
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
697 doc_ids: list[int] = list(
698 Document.objects.filter(tags=tag).values_list("pk", flat=True),
699 )
701 if not doc_ids:
702 return
704 parent_ids = [new_parent.id, *new_parent.get_ancestors_pks()]
706 parent_ids = list(dict.fromkeys(parent_ids))
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 )
715 to_create: list = []
716 affected: set[int] = set()
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
723 to_create.append(
724 doc_tag_relationship(document_id=doc_id, tag_id=parent_id),
725 )
726 affected.add(doc_id)
728 if to_create:
729 doc_tag_relationship.objects.bulk_create(
730 to_create,
731 ignore_conflicts=True,
732 )
734 if affected:
735 bulk_update_documents.apply_async(
736 kwargs={"document_ids": list(affected)},
737 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
738 )
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
752 from paperless_ai.indexing import update_llm_index
754 return update_llm_index(
755 iter_wrapper=iter_wrapper,
756 rebuild=rebuild,
757 )
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
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
787 if not apply_ai_suggestions_to_document(action, document):
788 return
790 # No document_updated signal to avoid loop
791 clear_document_caches(document.pk)
792 index_document.delay(document.pk)
794 ai_config = AIConfig()
795 if ai_config.llm_index_enabled:
796 update_document_in_llm_index.apply_async(kwargs={"document": document})
799@shared_task
800def update_document_in_llm_index(document) -> None:
801 llm_index_add_or_update_document(document)
804@shared_task
805def remove_document_from_llm_index(document) -> None:
806 llm_index_remove_document(document)
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
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 )
837 documents = list(bundle.documents.all().order_by("pk"))
839 _, temp_zip_path_str = mkstemp(suffix=".zip", dir=settings.SCRATCH_DIR)
840 temp_zip_path = Path(temp_zip_path_str)
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)
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)
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
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)