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

1from __future__ import annotations 

2 

3import logging 

4import tempfile 

5import uuid 

6from pathlib import Path 

7from typing import TYPE_CHECKING 

8from typing import Literal 

9from typing import NamedTuple 

10 

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 

19 

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 

39 

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 

42 

43 from django.contrib.auth.models import User 

44 

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 

47 

48logger: logging.Logger = logging.getLogger("paperless.bulk_edit") 

49 

50SourceMode = Literal["latest_version", "explicit_selection"] 

51 

52 

53class SourceModeChoices: 

54 LATEST_VERSION: SourceMode = "latest_version" 

55 EXPLICIT_SELECTION: SourceMode = "explicit_selection" 

56 

57 

58class ResolvedDocPair(NamedTuple): 

59 root_doc: Document 

60 source_doc: Document 

61 

62 

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) 

71 

72 

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 

87 

88 

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

97 

98 

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) 

105 

106 if source_mode == SourceModeChoices.EXPLICIT_SELECTION: 

107 return ResolvedDocPair(root_doc=root_doc, source_doc=doc) 

108 

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) 

112 

113 return ResolvedDocPair( 

114 root_doc=root_doc, 

115 source_doc=get_latest_version_for_root(root_doc), 

116 ) 

117 

118 

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) 

125 

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) 

133 

134 bulk_update_documents.apply_async( 

135 kwargs={"document_ids": affected_docs}, 

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

137 ) 

138 

139 return "OK" 

140 

141 

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) 

145 

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) 

155 

156 bulk_update_documents.apply_async( 

157 kwargs={"document_ids": affected_docs}, 

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

159 ) 

160 

161 return "OK" 

162 

163 

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) 

167 

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) 

175 

176 bulk_update_documents.apply_async( 

177 kwargs={"document_ids": affected_docs}, 

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

179 ) 

180 

181 return "OK" 

182 

183 

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

187 

188 DocumentTagRelationship = Document.tags.through 

189 to_create = [] 

190 affected_docs: set[int] = set() 

191 

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 ) 

200 

201 if to_create: 

202 DocumentTagRelationship.objects.bulk_create(to_create) 

203 

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 ) 

209 

210 return "OK" 

211 

212 

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

216 

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

224 

225 if affected_docs: 

226 bulk_update_documents.apply_async( 

227 kwargs={"document_ids": affected_docs}, 

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

229 ) 

230 

231 return "OK" 

232 

233 

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 

242 

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

249 

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

256 

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

263 

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 ) 

272 

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 ] 

279 

280 if to_create: 

281 DocumentTagRelationship.objects.bulk_create( 

282 to_create, 

283 ignore_conflicts=True, 

284 ) 

285 

286 if affected_docs: 

287 bulk_update_documents.apply_async( 

288 kwargs={"document_ids": affected_docs}, 

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

290 ) 

291 

292 return "OK" 

293 

294 

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 ) 

308 

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) 

339 

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 ) 

354 

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

360 

361 bulk_update_documents.apply_async( 

362 kwargs={"document_ids": affected_docs}, 

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

364 ) 

365 

366 return "OK" 

367 

368 

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

384 

385 Document.objects.filter(id__in=delete_ids).delete(transaction_id=uuid.uuid4()) 

386 

387 from documents.search import get_backend 

388 

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

390 for id in delete_ids: 

391 batch.remove(id) 

392 

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

401 

402 return "OK" 

403 

404 

405def reprocess(doc_ids: list[int], *, remote_ocr: bool = False) -> Literal["OK"]: 

406 """ 

407 Re-run parsing for the given documents. 

408 

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 ) 

417 

418 return "OK" 

419 

420 

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

429 

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) 

435 

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 ) 

443 

444 bulk_update_documents.apply_async( 

445 kwargs={"document_ids": affected_docs}, 

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

447 ) 

448 

449 return "OK" 

450 

451 

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) 

476 

477 import pikepdf 

478 

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) 

496 

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 

501 

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

518 

519 return "OK" 

520 

521 

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 

539 

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" 

570 

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

580 

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

594 

595 if user is not None: 

596 overrides.owner_id = user.id 

597 if not delete_originals: 

598 overrides.skip_asn_if_exists = True 

599 

600 if delete_originals and handoff_asn is not None: 

601 overrides.asn = handoff_asn 

602 

603 logger.info("Adding merged document to the task queue.") 

604 

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

612 

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

628 

629 return "OK" 

630 

631 

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} 

648 

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 ) 

657 

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 ] 

664 

665 updated_fields = ["root_document", "version_index", "archive_serial_number"] 

666 if version_label is not None: 

667 updated_fields.append("version_label") 

668 

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 

677 

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 ) 

683 

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 ) 

697 

698 Document.objects.filter(pk=root_document.pk).update(**root_updates) 

699 

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 ) 

712 

713 # One batch rather than a task each 

714 from documents.search import SearchIndexLockError 

715 from documents.search import get_backend 

716 

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) 

728 

729 bulk_update_documents.apply_async( 

730 kwargs={"document_ids": [root_document_id]}, 

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

732 ) 

733 

734 # And as far as the frontend is concerned, they're deleted 

735 status_mgr = DocumentsStatusManager() 

736 status_mgr.send_documents_deleted(source_ids) 

737 

738 return "OK" 

739 

740 

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 

756 

757 consume_tasks = [] 

758 

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

774 

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 ) 

795 

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

813 

814 except Exception as e: 

815 logger.exception(f"Error splitting document {doc.id}: {e}") 

816 

817 return "OK" 

818 

819 

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 

835 

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) 

849 

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

869 

870 return "OK" 

871 

872 

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

891 

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 

898 

899 pdf_docs: list[pikepdf.Pdf] = [] 

900 

901 try: 

902 if not operations: 

903 raise ValueError("Output document index is out of bounds") 

904 

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

911 

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

917 

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

921 

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) 

928 

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 ) 

987 

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

1002 

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 

1008 

1009 return "OK" 

1010 

1011 

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 

1028 

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 

1054 

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) 

1062 

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 

1094 

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 ) 

1104 

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

1112 

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 

1120 

1121 return "OK" 

1122 

1123 

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

1132 

1133 if target_doc_ids is None: 

1134 target_doc_ids = [] 

1135 

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 ) 

1149 

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) 

1178 

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

1185 

1186 

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