Coverage for documents/bulk_edit.py: 17%
514 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
1from __future__ import annotations
3import logging
4import tempfile
5import uuid
6from pathlib import Path
7from typing import TYPE_CHECKING
8from typing import Literal
9from typing import NamedTuple
11from celery import chord
12from celery import group
13from celery import shared_task
14from django.conf import settings
15from django.db import transaction
16from django.db.models import Max
17from django.db.models import Q
18from django.utils import timezone
20from documents.data_models import ConsumableDocument
21from documents.data_models import DocumentMetadataOverrides
22from documents.data_models import DocumentSource
23from documents.models import Correspondent
24from documents.models import CustomField
25from documents.models import CustomFieldInstance
26from documents.models import Document
27from documents.models import DocumentType
28from documents.models import PaperlessTask
29from documents.models import StoragePath
30from documents.models import Tag
31from documents.permissions import set_permissions_for_objects
32from documents.plugins.helpers import DocumentsStatusManager
33from documents.tasks import bulk_update_documents
34from documents.tasks import consume_file
35from documents.tasks import remove_document_from_index
36from documents.tasks import update_document_content_maybe_archive_file
37from documents.versioning import get_latest_version_for_root
38from documents.versioning import get_root_document
40if TYPE_CHECKING: 40 ↛ 41line 40 didn't jump to line 41 because the condition on line 40 was never true
41 from collections.abc import Mapping
43 from django.contrib.auth.models import User
45if settings.AUDIT_LOG_ENABLED: 45 ↛ 48line 45 didn't jump to line 48 because the condition on line 45 was always true
46 from auditlog.models import LogEntry
48logger: logging.Logger = logging.getLogger("paperless.bulk_edit")
50SourceMode = Literal["latest_version", "explicit_selection"]
53class SourceModeChoices:
54 LATEST_VERSION: SourceMode = "latest_version"
55 EXPLICIT_SELECTION: SourceMode = "explicit_selection"
58class ResolvedDocPair(NamedTuple):
59 root_doc: Document
60 source_doc: Document
63@shared_task(bind=True)
64def restore_archive_serial_numbers_task(
65 self,
66 backup: dict[int, int | None],
67 *args,
68 **kwargs,
69) -> None:
70 restore_archive_serial_numbers(backup)
73def release_archive_serial_numbers(doc_ids: list[int]) -> dict[int, int | None]:
74 """
75 Clears ASNs on documents that are about to be replaced so new documents
76 can be assigned ASNs without uniqueness collisions. Returns a backup map
77 of doc_id -> previous ASN for potential restoration.
78 """
79 qs = Document.objects.filter(
80 id__in=doc_ids,
81 archive_serial_number__isnull=False,
82 ).only("pk", "archive_serial_number")
83 backup = dict(qs.values_list("pk", "archive_serial_number"))
84 qs.update(archive_serial_number=None)
85 logger.info(f"Released archive serial numbers for documents {list(backup.keys())}")
86 return backup
89def restore_archive_serial_numbers(backup: dict[int, int | None]) -> None:
90 """
91 Restores ASNs using the provided backup map, intended for
92 rollback when replacement consumption fails.
93 """
94 for doc_id, asn in backup.items():
95 Document.objects.filter(pk=doc_id).update(archive_serial_number=asn)
96 logger.info(f"Restored archive serial numbers for documents {list(backup.keys())}")
99def _resolve_root_and_source_doc(
100 doc: Document,
101 *,
102 source_mode: SourceMode = SourceModeChoices.LATEST_VERSION,
103) -> ResolvedDocPair:
104 root_doc = get_root_document(doc)
106 if source_mode == SourceModeChoices.EXPLICIT_SELECTION:
107 return ResolvedDocPair(root_doc=root_doc, source_doc=doc)
109 # Version IDs are explicit by default, only a selected root resolves to latest
110 if doc.root_document_id is not None:
111 return ResolvedDocPair(root_doc=root_doc, source_doc=doc)
113 return ResolvedDocPair(
114 root_doc=root_doc,
115 source_doc=get_latest_version_for_root(root_doc),
116 )
119def set_correspondent(
120 doc_ids: list[int],
121 correspondent: Correspondent,
122) -> Literal["OK"]:
123 if correspondent:
124 correspondent = Correspondent.objects.only("pk").get(id=correspondent)
126 qs = (
127 Document.objects.filter(Q(id__in=doc_ids) & ~Q(correspondent=correspondent))
128 .select_related("correspondent")
129 .only("pk", "correspondent__id")
130 )
131 affected_docs = list(qs.values_list("pk", flat=True))
132 qs.update(correspondent=correspondent)
134 bulk_update_documents.apply_async(
135 kwargs={"document_ids": affected_docs},
136 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
137 )
139 return "OK"
142def set_storage_path(doc_ids: list[int], storage_path: StoragePath) -> Literal["OK"]:
143 if storage_path:
144 storage_path = StoragePath.objects.only("pk").get(id=storage_path)
146 qs = (
147 Document.objects.filter(
148 Q(id__in=doc_ids) & ~Q(storage_path=storage_path),
149 )
150 .select_related("storage_path")
151 .only("pk", "storage_path__id")
152 )
153 affected_docs = list(qs.values_list("pk", flat=True))
154 qs.update(storage_path=storage_path)
156 bulk_update_documents.apply_async(
157 kwargs={"document_ids": affected_docs},
158 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
159 )
161 return "OK"
164def set_document_type(doc_ids: list[int], document_type: DocumentType) -> Literal["OK"]:
165 if document_type:
166 document_type = DocumentType.objects.only("pk").get(id=document_type)
168 qs = (
169 Document.objects.filter(Q(id__in=doc_ids) & ~Q(document_type=document_type))
170 .select_related("document_type")
171 .only("pk", "document_type__id")
172 )
173 affected_docs = list(qs.values_list("pk", flat=True))
174 qs.update(document_type=document_type)
176 bulk_update_documents.apply_async(
177 kwargs={"document_ids": affected_docs},
178 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
179 )
181 return "OK"
184def add_tag(doc_ids: list[int], tag: int) -> Literal["OK"]:
185 tag_obj = Tag.objects.get(pk=tag)
186 tags_to_add = [tag_obj, *tag_obj.get_ancestors()]
188 DocumentTagRelationship = Document.tags.through
189 to_create = []
190 affected_docs: set[int] = set()
192 for t in tags_to_add:
193 qs = Document.objects.filter(Q(id__in=doc_ids) & ~Q(tags__id=t.id)).only("pk")
194 doc_ids_missing_tag = list(qs.values_list("pk", flat=True))
195 affected_docs.update(doc_ids_missing_tag)
196 to_create.extend(
197 DocumentTagRelationship(document_id=doc, tag_id=t.id)
198 for doc in doc_ids_missing_tag
199 )
201 if to_create:
202 DocumentTagRelationship.objects.bulk_create(to_create)
204 if affected_docs:
205 bulk_update_documents.apply_async(
206 kwargs={"document_ids": list(affected_docs)},
207 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
208 )
210 return "OK"
213def remove_tag(doc_ids: list[int], tag: int) -> Literal["OK"]:
214 tag_obj = Tag.objects.get(pk=tag)
215 tag_ids = [tag_obj.id, *tag_obj.get_descendants_pks()]
217 DocumentTagRelationship = Document.tags.through
218 qs = DocumentTagRelationship.objects.filter(
219 document_id__in=doc_ids,
220 tag_id__in=tag_ids,
221 )
222 affected_docs = list(qs.values_list("document_id", flat=True).distinct())
223 qs.delete()
225 if affected_docs:
226 bulk_update_documents.apply_async(
227 kwargs={"document_ids": affected_docs},
228 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
229 )
231 return "OK"
234def modify_tags(
235 doc_ids: list[int],
236 add_tags: list[int],
237 remove_tags: list[int],
238) -> Literal["OK"]:
239 qs = Document.objects.filter(id__in=doc_ids).only("pk")
240 affected_docs = list(qs.values_list("pk", flat=True))
241 DocumentTagRelationship = Document.tags.through
243 # add with all ancestors
244 expanded_add_tags: set[int] = set()
245 add_tag_objects = Tag.objects.filter(pk__in=add_tags)
246 for t in add_tag_objects:
247 expanded_add_tags.add(int(t.id))
248 expanded_add_tags.update(int(pk) for pk in t.get_ancestors_pks())
250 # remove with all descendants
251 expanded_remove_tags: set[int] = set()
252 remove_tag_objects = Tag.objects.filter(pk__in=remove_tags)
253 for t in remove_tag_objects:
254 expanded_remove_tags.add(int(t.id))
255 expanded_remove_tags.update(int(pk) for pk in t.get_descendants_pks())
257 with transaction.atomic():
258 if expanded_remove_tags:
259 DocumentTagRelationship.objects.filter(
260 document_id__in=affected_docs,
261 tag_id__in=expanded_remove_tags,
262 ).delete()
264 to_create = []
265 if expanded_add_tags:
266 existing_pairs = set(
267 DocumentTagRelationship.objects.filter(
268 document_id__in=affected_docs,
269 tag_id__in=expanded_add_tags,
270 ).values_list("document_id", "tag_id"),
271 )
273 to_create = [
274 DocumentTagRelationship(document_id=doc, tag_id=tag)
275 for doc in affected_docs
276 for tag in expanded_add_tags
277 if (doc, tag) not in existing_pairs
278 ]
280 if to_create:
281 DocumentTagRelationship.objects.bulk_create(
282 to_create,
283 ignore_conflicts=True,
284 )
286 if affected_docs:
287 bulk_update_documents.apply_async(
288 kwargs={"document_ids": affected_docs},
289 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
290 )
292 return "OK"
295def modify_custom_fields(
296 doc_ids: list[int],
297 add_custom_fields: list[int] | dict,
298 remove_custom_fields: list[int],
299) -> Literal["OK"]:
300 qs = Document.objects.filter(id__in=doc_ids).only("pk")
301 affected_docs = list(qs.values_list("pk", flat=True))
302 # Ensure add_custom_fields is a list of (int, value) tuples, supports old API
303 add_custom_fields = (
304 [(int(field), value) for field, value in add_custom_fields.items()]
305 if isinstance(add_custom_fields, dict)
306 else [(int(field), None) for field in add_custom_fields]
307 )
309 # Resolved once, instead of re-querying the same field for every document
310 custom_fields_by_id: dict[int, CustomField] = CustomField.objects.in_bulk(
311 [field_id for field_id, _ in add_custom_fields],
312 )
313 # Passed to update_or_create() below rather than a bare id, so the FK is
314 # cached on the created instance and auditlog's post_save receiver does
315 # not reload it per row. Only needed for additions. content is deferred:
316 # the one field here that is both large and unused.
317 docs_by_id: dict[int, Document] = (
318 Document.objects.defer("content").in_bulk(affected_docs)
319 if add_custom_fields
320 else {}
321 )
322 for field_id, value in add_custom_fields:
323 custom_field = custom_fields_by_id[field_id]
324 value_field = CustomFieldInstance.TYPE_TO_DATA_STORE_NAME_MAP[
325 custom_field.data_type
326 ]
327 is_doclink = custom_field.data_type == CustomField.FieldDataType.DOCUMENTLINK
328 for doc_id in affected_docs:
329 if is_doclink and value and doc_id in value:
330 # Prevent self-linking
331 continue
332 CustomFieldInstance.objects.update_or_create(
333 document=docs_by_id[doc_id],
334 field=custom_field,
335 defaults={value_field: value},
336 )
337 if is_doclink:
338 reflect_doclinks(docs_by_id[doc_id], custom_field, value)
340 # For doc link fields that are being removed, remove symmetrical links.
341 # select_related avoids a per-instance reload of the document and field.
342 for doclink_being_removed_instance in CustomFieldInstance.objects.filter(
343 document_id__in=affected_docs,
344 field__id__in=remove_custom_fields,
345 field__data_type=CustomField.FieldDataType.DOCUMENTLINK,
346 value_document_ids__isnull=False,
347 ).select_related("field", "document"):
348 for target_doc_id in doclink_being_removed_instance.value:
349 remove_doclink(
350 document=doclink_being_removed_instance.document,
351 field=doclink_being_removed_instance.field,
352 target_doc_id=target_doc_id,
353 )
355 # Finally, remove the custom fields
356 CustomFieldInstance.objects.filter(
357 document_id__in=affected_docs,
358 field_id__in=remove_custom_fields,
359 ).hard_delete()
361 bulk_update_documents.apply_async(
362 kwargs={"document_ids": affected_docs},
363 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
364 )
366 return "OK"
369@shared_task
370def delete(doc_ids: list[int]) -> Literal["OK"]:
371 try:
372 root_ids = (
373 Document.objects.filter(id__in=doc_ids, root_document__isnull=True)
374 .values_list("id", flat=True)
375 .distinct()
376 )
377 version_ids = (
378 Document.objects.filter(root_document_id__in=root_ids)
379 .exclude(id__in=doc_ids)
380 .values_list("id", flat=True)
381 .distinct()
382 )
383 delete_ids = list({*doc_ids, *version_ids})
385 Document.objects.filter(id__in=delete_ids).delete(transaction_id=uuid.uuid4())
387 from documents.search import get_backend
389 with get_backend().batch_update() as batch:
390 for id in delete_ids:
391 batch.remove(id)
393 status_mgr = DocumentsStatusManager()
394 status_mgr.send_documents_deleted(delete_ids)
395 except Exception as e:
396 if "Data too long for column" in str(e):
397 logger.warning(
398 "Detected a possible incompatible database column. See https://docs.paperless-ngx.com/troubleshooting/#convert-uuid-field",
399 )
400 logger.error(f"Error deleting documents: {e!s}")
402 return "OK"
405def reprocess(doc_ids: list[int], *, remote_ocr: bool = False) -> Literal["OK"]:
406 """
407 Re-run parsing for the given documents.
409 Consumption workflows do not run here, so ``remote_ocr`` is how the user
410 asks for the remote engine when it is not configured to handle everything.
411 """
412 for document_id in doc_ids: 412 ↛ 413line 412 didn't jump to line 413 because the loop on line 412 never started
413 update_document_content_maybe_archive_file.apply_async(
414 kwargs={"document_id": document_id, "remote_ocr": remote_ocr},
415 headers={"trigger_source": PaperlessTask.TriggerSource.MANUAL},
416 )
418 return "OK"
421def set_permissions(
422 doc_ids: list[int],
423 set_permissions: dict,
424 *,
425 owner: User | None = None,
426 merge: bool = False,
427) -> Literal["OK"]:
428 qs = Document.objects.filter(id__in=doc_ids).select_related("owner")
430 if merge:
431 # If merging, only set owner for documents that don't have an owner
432 qs.filter(owner__isnull=True).update(owner=owner)
433 else:
434 qs.update(owner=owner)
436 affected_docs = list(qs.values_list("pk", flat=True))
437 set_permissions_for_objects(
438 permissions=set_permissions,
439 model=Document,
440 pks=affected_docs,
441 merge=merge,
442 )
444 bulk_update_documents.apply_async(
445 kwargs={"document_ids": affected_docs},
446 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
447 )
449 return "OK"
452def rotate(
453 doc_ids: list[int],
454 degrees: int,
455 *,
456 source_mode: SourceMode = SourceModeChoices.LATEST_VERSION,
457 user: User | None = None,
458 trigger_source: PaperlessTask.TriggerSource = PaperlessTask.TriggerSource.WEB_UI,
459) -> Literal["OK"]:
460 logger.info(
461 f"Attempting to rotate {len(doc_ids)} documents by {degrees} degrees.",
462 )
463 docs_by_id = {
464 doc.id: doc
465 for doc in Document.objects.select_related("root_document").filter(
466 id__in=doc_ids,
467 )
468 }
469 docs_by_root_id: dict[int, ResolvedDocPair] = {}
470 for doc_id in doc_ids: 470 ↛ 471line 470 didn't jump to line 471 because the loop on line 470 never started
471 doc = docs_by_id.get(doc_id)
472 if doc is None:
473 continue
474 pair = _resolve_root_and_source_doc(doc, source_mode=source_mode)
475 docs_by_root_id.setdefault(pair.root_doc.id, pair)
477 import pikepdf
479 for pair in docs_by_root_id.values(): 479 ↛ 480line 479 didn't jump to line 480 because the loop on line 479 never started
480 if pair.source_doc.mime_type != "application/pdf":
481 logger.warning(
482 f"Document {pair.root_doc.id} is not a PDF, skipping rotation.",
483 )
484 continue
485 try:
486 # Write rotated output to a temp file and create a new version via consume pipeline
487 filepath: Path = (
488 Path(tempfile.mkdtemp(dir=settings.SCRATCH_DIR))
489 / f"{pair.root_doc.id}_rotated.pdf"
490 )
491 with pikepdf.open(pair.source_doc.source_path) as pdf:
492 for page in pdf.pages:
493 page.rotate(degrees, relative=True)
494 pdf.remove_unreferenced_resources()
495 pdf.save(filepath)
497 # Preserve metadata/permissions via overrides; mark as new version
498 overrides = DocumentMetadataOverrides().from_document(pair.root_doc)
499 if user is not None:
500 overrides.actor_id = user.id
502 consume_file.apply_async(
503 kwargs={
504 "input_doc": ConsumableDocument(
505 source=DocumentSource.ConsumeFolder,
506 original_file=filepath,
507 root_document_id=pair.root_doc.id,
508 ),
509 "overrides": overrides,
510 },
511 headers={"trigger_source": trigger_source},
512 )
513 logger.info(
514 f"Queued new rotated version for document {pair.root_doc.id} by {degrees} degrees",
515 )
516 except Exception as e:
517 logger.exception(f"Error rotating document {pair.root_doc.id}: {e}")
519 return "OK"
522def merge(
523 doc_ids: list[int],
524 *,
525 metadata_document_id: int | None = None,
526 delete_originals: bool = False,
527 archive_fallback: bool = False,
528 source_mode: SourceMode = SourceModeChoices.LATEST_VERSION,
529 user: User | None = None,
530 trigger_source: PaperlessTask.TriggerSource = PaperlessTask.TriggerSource.WEB_UI,
531) -> Literal["OK"]:
532 logger.info(
533 f"Attempting to merge {len(doc_ids)} documents into a single document.",
534 )
535 qs = Document.objects.select_related("root_document").filter(id__in=doc_ids)
536 docs_by_id = {doc.id: doc for doc in qs}
537 affected_docs: list[int] = []
538 import pikepdf
540 merged_pdf = pikepdf.new()
541 version: str = merged_pdf.pdf_version
542 handoff_asn: int | None = None
543 # use doc_ids to preserve order
544 for doc_id in doc_ids: 544 ↛ 545line 544 didn't jump to line 545 because the loop on line 544 never started
545 doc = docs_by_id.get(doc_id)
546 if doc is None:
547 continue
548 pair = _resolve_root_and_source_doc(doc, source_mode=source_mode)
549 try:
550 doc_path = (
551 pair.source_doc.archive_path
552 if archive_fallback
553 and pair.source_doc.mime_type != "application/pdf"
554 and pair.source_doc.has_archive_version
555 else pair.source_doc.source_path
556 )
557 with pikepdf.open(str(doc_path)) as pdf:
558 version = max(version, pdf.pdf_version)
559 merged_pdf.pages.extend(pdf.pages)
560 affected_docs.append(doc.id)
561 if handoff_asn is None and doc.archive_serial_number is not None:
562 handoff_asn = doc.archive_serial_number
563 except Exception as e:
564 logger.exception(
565 f"Error merging document {doc.id}, it will not be included in the merge: {e}",
566 )
567 if len(affected_docs) == 0: 567 ↛ 571line 567 didn't jump to line 571 because the condition on line 567 was always true
568 logger.warning("No documents were merged")
569 return "OK"
571 filepath = (
572 Path(
573 tempfile.mkdtemp(dir=settings.SCRATCH_DIR),
574 )
575 / f"{'_'.join([str(doc_id) for doc_id in affected_docs])[:100]}_merged.pdf"
576 )
577 merged_pdf.remove_unreferenced_resources()
578 merged_pdf.save(filepath, min_version=version)
579 merged_pdf.close()
581 if metadata_document_id:
582 metadata_document = qs.get(id=metadata_document_id)
583 if metadata_document is not None:
584 overrides: DocumentMetadataOverrides = (
585 DocumentMetadataOverrides.from_document(metadata_document)
586 )
587 overrides.title = metadata_document.title + " (merged)"
588 if metadata_document.archive_serial_number is not None:
589 handoff_asn = metadata_document.archive_serial_number
590 else:
591 overrides = DocumentMetadataOverrides()
592 else:
593 overrides = DocumentMetadataOverrides()
595 if user is not None:
596 overrides.owner_id = user.id
597 if not delete_originals:
598 overrides.skip_asn_if_exists = True
600 if delete_originals and handoff_asn is not None:
601 overrides.asn = handoff_asn
603 logger.info("Adding merged document to the task queue.")
605 consume_task = consume_file.s(
606 input_doc=ConsumableDocument(
607 source=DocumentSource.ConsumeFolder,
608 original_file=filepath,
609 ),
610 overrides=overrides,
611 ).set(headers={"trigger_source": trigger_source})
613 if delete_originals:
614 backup = release_archive_serial_numbers(affected_docs)
615 logger.info(
616 "Queueing removal of original documents after consumption of merged document",
617 )
618 try:
619 consume_task.apply_async(
620 link=[delete.si(affected_docs)],
621 link_error=[restore_archive_serial_numbers_task.s(backup)],
622 )
623 except Exception:
624 restore_archive_serial_numbers(backup)
625 raise
626 else:
627 consume_task.apply_async()
629 return "OK"
632def merge_as_versions(
633 doc_ids: list[int],
634 *,
635 root_document_id: int,
636 version_label: str | None = None,
637 user: User | None = None,
638) -> Literal["OK"]:
639 with transaction.atomic():
640 documents = list(
641 # Ordered by pk so concurrent merges take the row locks in the same order
642 Document.objects.select_for_update()
643 .filter(id__in=doc_ids)
644 .order_by("id")
645 .defer("content"),
646 )
647 documents_by_id = {document.id: document for document in documents}
649 source_ids = [doc_id for doc_id in doc_ids if doc_id != root_document_id]
650 root_document = documents_by_id[root_document_id]
651 next_version_index = (
652 Document.global_objects.filter(
653 root_document_id=root_document_id,
654 ).aggregate(max_index=Max("version_index"))["max_index"]
655 or 0
656 )
658 # A version gives up its ASN
659 source_asns = [
660 documents_by_id[source_id].archive_serial_number
661 for source_id in source_ids
662 if documents_by_id[source_id].archive_serial_number is not None
663 ]
665 updated_fields = ["root_document", "version_index", "archive_serial_number"]
666 if version_label is not None:
667 updated_fields.append("version_label")
669 for source_id in source_ids:
670 next_version_index += 1
671 source_document = documents_by_id[source_id]
672 source_document.root_document_id = root_document.pk
673 source_document.version_index = next_version_index
674 source_document.archive_serial_number = None
675 if version_label is not None:
676 source_document.version_label = version_label
678 # bulk_update and not save() to avoid post_save now
679 Document.objects.bulk_update(
680 [documents_by_id[source_id] for source_id in source_ids],
681 updated_fields,
682 )
684 root_updates = {"modified": timezone.now()}
685 if source_asns and root_document.archive_serial_number is None:
686 # If a version had one, hand the ASN over, the same as merge() does
687 root_updates["archive_serial_number"] = source_asns.pop(0)
688 logger.info(
689 f"Document {root_document.id} took archive serial number "
690 f"{root_updates['archive_serial_number']} from a document merged into it",
691 )
692 if source_asns:
693 logger.warning(
694 f"Archive serial number(s) {source_asns} were removed by merging "
695 f"those documents as versions of document {root_document.id}",
696 )
698 Document.objects.filter(pk=root_document.pk).update(**root_updates)
700 if settings.AUDIT_LOG_ENABLED:
701 # update() doesn't fire auditlog signals, so manual
702 LogEntry.objects.log_create(
703 instance=root_document,
704 changes={"Merged As Versions": ["None", source_ids]},
705 action=LogEntry.Action.UPDATE,
706 actor=user,
707 additional_data={
708 "reason": "Merged as versions",
709 "version_ids": source_ids,
710 },
711 )
713 # One batch rather than a task each
714 from documents.search import SearchIndexLockError
715 from documents.search import get_backend
717 try:
718 with get_backend().batch_update() as batch:
719 for source_id in source_ids:
720 batch.remove(source_id)
721 except SearchIndexLockError:
722 logger.error(
723 f"Search index lock exhausted removing {source_ids}, "
724 f"scheduling deferred index removal",
725 )
726 for source_id in source_ids:
727 remove_document_from_index.apply_async(args=[source_id], countdown=60)
729 bulk_update_documents.apply_async(
730 kwargs={"document_ids": [root_document_id]},
731 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM},
732 )
734 # And as far as the frontend is concerned, they're deleted
735 status_mgr = DocumentsStatusManager()
736 status_mgr.send_documents_deleted(source_ids)
738 return "OK"
741def split(
742 doc_ids: list[int],
743 pages: list[list[int]],
744 *,
745 delete_originals: bool = False,
746 source_mode: SourceMode = SourceModeChoices.LATEST_VERSION,
747 user: User | None = None,
748 trigger_source: PaperlessTask.TriggerSource = PaperlessTask.TriggerSource.WEB_UI,
749) -> Literal["OK"]:
750 logger.info(
751 f"Attempting to split document {doc_ids[0]} into {len(pages)} documents",
752 )
753 doc = Document.objects.select_related("root_document").get(id=doc_ids[0])
754 pair = _resolve_root_and_source_doc(doc, source_mode=source_mode)
755 import pikepdf
757 consume_tasks = []
759 try:
760 with pikepdf.open(pair.source_doc.source_path) as pdf:
761 for idx, split_doc in enumerate(pages):
762 dst: pikepdf.Pdf = pikepdf.new()
763 for page in split_doc:
764 dst.pages.append(pdf.pages[page - 1])
765 filepath: Path = (
766 Path(
767 tempfile.mkdtemp(dir=settings.SCRATCH_DIR),
768 )
769 / f"{doc.id}_{split_doc[0]}-{split_doc[-1]}.pdf"
770 )
771 dst.remove_unreferenced_resources()
772 dst.save(filepath)
773 dst.close()
775 overrides: DocumentMetadataOverrides = (
776 DocumentMetadataOverrides().from_document(doc)
777 )
778 overrides.title = f"{doc.title} (split {idx + 1})"
779 if user is not None:
780 overrides.owner_id = user.id
781 if not delete_originals:
782 overrides.skip_asn_if_exists = True
783 logger.info(
784 f"Adding split document with pages {split_doc} to the task queue.",
785 )
786 consume_tasks.append(
787 consume_file.s(
788 input_doc=ConsumableDocument(
789 source=DocumentSource.ConsumeFolder,
790 original_file=filepath,
791 ),
792 overrides=overrides,
793 ).set(headers={"trigger_source": trigger_source}),
794 )
796 if delete_originals:
797 backup = release_archive_serial_numbers([doc.id])
798 logger.info(
799 "Queueing removal of original document after consumption of the split documents",
800 )
801 try:
802 chord(
803 header=consume_tasks,
804 body=delete.si([doc.id]),
805 ).on_error(
806 restore_archive_serial_numbers_task.s(backup),
807 ).apply_async()
808 except Exception:
809 restore_archive_serial_numbers(backup)
810 raise
811 else:
812 group(consume_tasks).delay()
814 except Exception as e:
815 logger.exception(f"Error splitting document {doc.id}: {e}")
817 return "OK"
820def delete_pages(
821 doc_ids: list[int],
822 pages: list[int],
823 *,
824 source_mode: SourceMode = SourceModeChoices.LATEST_VERSION,
825 user: User | None = None,
826 trigger_source: PaperlessTask.TriggerSource = PaperlessTask.TriggerSource.WEB_UI,
827) -> Literal["OK"]:
828 logger.info(
829 f"Attempting to delete pages {pages} from {len(doc_ids)} documents",
830 )
831 doc = Document.objects.select_related("root_document").get(id=doc_ids[0])
832 pair = _resolve_root_and_source_doc(doc, source_mode=source_mode)
833 pages = sorted(pages) # sort pages to avoid index issues
834 import pikepdf
836 try:
837 # Produce edited PDF to a temp file and create a new version
838 filepath: Path = (
839 Path(tempfile.mkdtemp(dir=settings.SCRATCH_DIR))
840 / f"{pair.root_doc.id}_pages_deleted.pdf"
841 )
842 with pikepdf.open(pair.source_doc.source_path) as pdf:
843 offset = 1 # pages are 1-indexed
844 for page_num in pages:
845 pdf.pages.remove(pdf.pages[page_num - offset])
846 offset += 1 # remove() changes the index of the pages
847 pdf.remove_unreferenced_resources()
848 pdf.save(filepath)
850 overrides = DocumentMetadataOverrides().from_document(pair.root_doc)
851 if user is not None:
852 overrides.actor_id = user.id
853 consume_file.apply_async(
854 kwargs={
855 "input_doc": ConsumableDocument(
856 source=DocumentSource.ConsumeFolder,
857 original_file=filepath,
858 root_document_id=pair.root_doc.id,
859 ),
860 "overrides": overrides,
861 },
862 headers={"trigger_source": trigger_source},
863 )
864 logger.info(
865 f"Queued new version for document {pair.root_doc.id} after deleting pages {pages}",
866 )
867 except Exception as e:
868 logger.exception(f"Error deleting pages from document {pair.root_doc.id}: {e}")
870 return "OK"
873def edit_pdf(
874 doc_ids: list[int],
875 operations: list[dict[str, int]],
876 *,
877 delete_original: bool = False,
878 update_document: bool = False,
879 include_metadata: bool = True,
880 source_mode: SourceMode = SourceModeChoices.LATEST_VERSION,
881 user: User | None = None,
882 trigger_source: PaperlessTask.TriggerSource = PaperlessTask.TriggerSource.WEB_UI,
883) -> Literal["OK"]:
884 """
885 Operations is a list of dictionaries describing the final PDF pages.
886 Each entry must contain the original page number in `page` and may
887 specify `rotate` in degrees and `doc` indicating the output
888 document index (for splitting). Pages omitted from the list are
889 discarded.
890 """
892 logger.info(
893 f"Editing PDF of document {doc_ids[0]} with {len(operations)} operations",
894 )
895 doc = Document.objects.select_related("root_document").get(id=doc_ids[0])
896 pair = _resolve_root_and_source_doc(doc, source_mode=source_mode)
897 import pikepdf
899 pdf_docs: list[pikepdf.Pdf] = []
901 try:
902 if not operations:
903 raise ValueError("Output document index is out of bounds")
905 max_idx = max(op.get("doc", 0) for op in operations)
906 if update_document and max_idx > 0:
907 logger.error(
908 "Update requested but multiple output documents specified",
909 )
910 raise ValueError("Multiple output documents specified")
912 if any(
913 op.get("doc", 0) < 0 or op.get("doc", 0) >= len(operations)
914 for op in operations
915 ):
916 raise ValueError("Output document index is out of bounds")
918 with pikepdf.open(pair.source_doc.source_path) as src:
919 # prepare output documents
920 pdf_docs = [pikepdf.new() for _ in range(max_idx + 1)]
922 for op in operations:
923 dst = pdf_docs[op.get("doc", 0)]
924 page = src.pages[op["page"] - 1]
925 dst.pages.append(page)
926 if op.get("rotate"):
927 dst.pages[-1].rotate(op["rotate"], relative=True)
929 if update_document:
930 # Create a new version from the edited PDF rather than replacing in-place
931 pdf = pdf_docs[0]
932 pdf.remove_unreferenced_resources()
933 filepath: Path = (
934 Path(tempfile.mkdtemp(dir=settings.SCRATCH_DIR))
935 / f"{pair.root_doc.id}_edited.pdf"
936 )
937 pdf.save(filepath)
938 overrides = (
939 DocumentMetadataOverrides().from_document(pair.root_doc)
940 if include_metadata
941 else DocumentMetadataOverrides()
942 )
943 if user is not None:
944 overrides.owner_id = user.id
945 overrides.actor_id = user.id
946 consume_file.apply_async(
947 kwargs={
948 "input_doc": ConsumableDocument(
949 source=DocumentSource.ConsumeFolder,
950 original_file=filepath,
951 root_document_id=pair.root_doc.id,
952 ),
953 "overrides": overrides,
954 },
955 headers={"trigger_source": trigger_source},
956 )
957 else:
958 consume_tasks = []
959 overrides = (
960 DocumentMetadataOverrides().from_document(pair.root_doc)
961 if include_metadata
962 else DocumentMetadataOverrides()
963 )
964 if user is not None:
965 overrides.owner_id = user.id
966 overrides.actor_id = user.id
967 if not delete_original:
968 overrides.skip_asn_if_exists = True
969 if delete_original and len(pdf_docs) == 1:
970 overrides.asn = pair.root_doc.archive_serial_number
971 for idx, pdf in enumerate(pdf_docs, start=1):
972 version_filepath: Path = (
973 Path(tempfile.mkdtemp(dir=settings.SCRATCH_DIR))
974 / f"{pair.root_doc.id}_edit_{idx}.pdf"
975 )
976 pdf.remove_unreferenced_resources()
977 pdf.save(version_filepath)
978 consume_tasks.append(
979 consume_file.s(
980 input_doc=ConsumableDocument(
981 source=DocumentSource.ConsumeFolder,
982 original_file=version_filepath,
983 ),
984 overrides=overrides,
985 ).set(headers={"trigger_source": trigger_source}),
986 )
988 if delete_original:
989 backup = release_archive_serial_numbers([doc.id])
990 try:
991 chord(
992 header=consume_tasks,
993 body=delete.si([doc.id]),
994 ).on_error(
995 restore_archive_serial_numbers_task.s(backup),
996 ).apply_async()
997 except Exception:
998 restore_archive_serial_numbers(backup)
999 raise
1000 else:
1001 group(consume_tasks).delay()
1003 except Exception as e:
1004 logger.exception(f"Error editing document {pair.root_doc.id}: {e}")
1005 raise ValueError(
1006 f"An error occurred while editing the document: {e}",
1007 ) from e
1009 return "OK"
1012def remove_password(
1013 doc_ids: list[int],
1014 password: str,
1015 *,
1016 update_document: bool = False,
1017 delete_original: bool = False,
1018 include_metadata: bool = True,
1019 source_mode: SourceMode = SourceModeChoices.LATEST_VERSION,
1020 user: User | None = None,
1021 trigger_source: PaperlessTask.TriggerSource = PaperlessTask.TriggerSource.WEB_UI,
1022 source_paths_by_id: Mapping[int, Path] | None = None,
1023) -> Literal["OK"]:
1024 """
1025 Remove password protection from PDF documents.
1026 """
1027 import pikepdf
1029 for doc_id in doc_ids: 1029 ↛ 1030line 1029 didn't jump to line 1030 because the loop on line 1029 never started
1030 doc = Document.objects.select_related("root_document").get(id=doc_id)
1031 pair = _resolve_root_and_source_doc(doc, source_mode=source_mode)
1032 try:
1033 logger.info(
1034 f"Attempting password removal from document {pair.root_doc.id}",
1035 )
1036 # The caller may supply an explicit source path (e.g. the staged
1037 # file during consumption, before source_path is populated).
1038 source_path = (source_paths_by_id or {}).get(
1039 doc.id,
1040 pair.source_doc.source_path,
1041 )
1042 try:
1043 with pikepdf.open(source_path) as pdf:
1044 if not pdf.is_encrypted:
1045 logger.info(
1046 "Skipping password removal for document %s because the "
1047 "source PDF is not encrypted",
1048 pair.root_doc.id,
1049 )
1050 continue
1051 except pikepdf.PasswordError:
1052 # Password-protected PDFs need the supplied password below.
1053 pass
1055 with pikepdf.open(source_path, password=password) as pdf:
1056 filepath: Path = (
1057 Path(tempfile.mkdtemp(dir=settings.SCRATCH_DIR))
1058 / f"{pair.root_doc.id}_unprotected.pdf"
1059 )
1060 pdf.remove_unreferenced_resources()
1061 pdf.save(filepath)
1063 if update_document:
1064 # Create a new version rather than modifying the root/original in place.
1065 overrides = (
1066 DocumentMetadataOverrides().from_document(pair.root_doc)
1067 if include_metadata
1068 else DocumentMetadataOverrides()
1069 )
1070 if user is not None:
1071 overrides.owner_id = user.id
1072 overrides.actor_id = user.id
1073 consume_file.apply_async(
1074 kwargs={
1075 "input_doc": ConsumableDocument(
1076 source=DocumentSource.ConsumeFolder,
1077 original_file=filepath,
1078 root_document_id=pair.root_doc.id,
1079 ),
1080 "overrides": overrides,
1081 },
1082 headers={"trigger_source": trigger_source},
1083 )
1084 else:
1085 consume_tasks = []
1086 overrides = (
1087 DocumentMetadataOverrides().from_document(pair.root_doc)
1088 if include_metadata
1089 else DocumentMetadataOverrides()
1090 )
1091 if user is not None:
1092 overrides.owner_id = user.id
1093 overrides.actor_id = user.id
1095 consume_tasks.append(
1096 consume_file.s(
1097 input_doc=ConsumableDocument(
1098 source=DocumentSource.ConsumeFolder,
1099 original_file=filepath,
1100 ),
1101 overrides=overrides,
1102 ).set(headers={"trigger_source": trigger_source}),
1103 )
1105 if delete_original:
1106 chord(
1107 header=consume_tasks,
1108 body=delete.si([doc.id]),
1109 ).delay()
1110 else:
1111 group(consume_tasks).delay()
1113 except Exception as e:
1114 logger.exception(
1115 f"Error removing password from document {pair.root_doc.id}: {e}",
1116 )
1117 raise ValueError(
1118 f"An error occurred while removing the password: {e}",
1119 ) from e
1121 return "OK"
1124def reflect_doclinks(
1125 document: Document,
1126 field: CustomField,
1127 target_doc_ids: list[int],
1128) -> None:
1129 """
1130 Add or remove 'symmetrical' links to `document` on all `target_doc_ids`
1131 """
1133 if target_doc_ids is None:
1134 target_doc_ids = []
1136 # Check if any documents are going to be removed from the current list of links and remove the symmetrical links
1137 current_field_instance = CustomFieldInstance.objects.filter(
1138 field=field,
1139 document=document,
1140 ).first()
1141 if current_field_instance is not None and current_field_instance.value is not None:
1142 for doc_id in current_field_instance.value:
1143 if doc_id not in target_doc_ids:
1144 remove_doclink(
1145 document=document,
1146 field=field,
1147 target_doc_id=doc_id,
1148 )
1150 # Create an instance if target doc doesn't have this field or append it to an existing one
1151 existing_custom_field_instances = {
1152 custom_field.document_id: custom_field
1153 for custom_field in CustomFieldInstance.objects.filter(
1154 field=field,
1155 document_id__in=target_doc_ids,
1156 )
1157 }
1158 custom_field_instances_to_create = []
1159 custom_field_instances_to_update = []
1160 for target_doc_id in target_doc_ids:
1161 target_doc_field_instance = existing_custom_field_instances.get(
1162 target_doc_id,
1163 )
1164 if target_doc_field_instance is None:
1165 custom_field_instances_to_create.append(
1166 CustomFieldInstance(
1167 document_id=target_doc_id,
1168 field=field,
1169 value_document_ids=[document.id],
1170 ),
1171 )
1172 elif target_doc_field_instance.value is None:
1173 target_doc_field_instance.value_document_ids = [document.id]
1174 custom_field_instances_to_update.append(target_doc_field_instance)
1175 elif document.id not in target_doc_field_instance.value:
1176 target_doc_field_instance.value_document_ids.append(document.id)
1177 custom_field_instances_to_update.append(target_doc_field_instance)
1179 CustomFieldInstance.objects.bulk_create(custom_field_instances_to_create)
1180 CustomFieldInstance.objects.bulk_update(
1181 custom_field_instances_to_update,
1182 ["value_document_ids"],
1183 )
1184 Document.objects.filter(id__in=target_doc_ids).update(modified=timezone.now())
1187def remove_doclink(
1188 document: Document,
1189 field: CustomField,
1190 target_doc_id: int,
1191) -> None:
1192 """
1193 Removes a 'symmetrical' link to `document` from the target document's existing custom field instance
1194 """
1195 # select_related: a signal receiver (auditlog) touches .document/.field on
1196 # the save() below, without this that is a per-call reload query
1197 target_doc_field_instance = (
1198 CustomFieldInstance.objects.filter(document_id=target_doc_id, field=field)
1199 .select_related("document", "field")
1200 .first()
1201 )
1202 if (
1203 target_doc_field_instance is not None
1204 and document.id in target_doc_field_instance.value
1205 ):
1206 target_doc_field_instance.value.remove(document.id)
1207 target_doc_field_instance.save()
1208 Document.objects.filter(id=target_doc_id).update(modified=timezone.now())