Compare commits

...
Author SHA1 Message Date
shamoonandTrenton H 67cd764ab5 Resolve reprocess with stumyplogs deduplicated query
Co-Authored-By: Trenton H <797416+stumpylog@users.noreply.github.com>
2026-10-08 09:57:11 -07:00
shamoon 745ef52fc8 Use DocumentFactory and mocker in reprocess tests 2026-10-08 09:55:24 -07:00
shamoon df02e6073f Cover LLM index update when reprocessing a version 2026-10-08 09:55:23 -07:00
shamoon 702c5561d4 Fix: bulk reprocess latest version for root documents 2026-10-07 16:14:36 -07:00
Trenton H d7a9894400 Fix: retry a search index rebuild that was interrupted (#14379)
An interrupted rebuild left an empty index stamped as current, so the next
start reported it as up to date. Mark the rebuild as in progress and only
clear the marker once it completes.
2026-10-07 08:54:36 -07:00
7 changed files with 336 additions and 68 deletions

No files matched your search

+27 -2
View File
@@ -13,8 +13,14 @@ from celery import group
from celery import shared_task from celery import shared_task
from django.conf import settings from django.conf import settings
from django.db import transaction from django.db import transaction
from django.db.models import Case
from django.db.models import F
from django.db.models import Max from django.db.models import Max
from django.db.models import OuterRef
from django.db.models import Q from django.db.models import Q
from django.db.models import Subquery
from django.db.models import When
from django.db.models.functions import Coalesce
from django.utils import timezone from django.utils import timezone
from documents.data_models import ConsumableDocument from documents.data_models import ConsumableDocument
@@ -36,6 +42,7 @@ from documents.tasks import remove_document_from_index
from documents.tasks import update_document_content_maybe_archive_file from documents.tasks import update_document_content_maybe_archive_file
from documents.versioning import get_latest_version_for_root from documents.versioning import get_latest_version_for_root
from documents.versioning import get_root_document from documents.versioning import get_root_document
from documents.versioning import versions_newest_first
if TYPE_CHECKING: if TYPE_CHECKING:
from collections.abc import Mapping from collections.abc import Mapping
@@ -408,10 +415,28 @@ def reprocess(doc_ids: list[int], *, remote_ocr: bool = False) -> Literal["OK"]:
Consumption workflows do not run here, so ``remote_ocr`` is how the user Consumption workflows do not run here, so ``remote_ocr`` is how the user
asks for the remote engine when it is not configured to handle everything. asks for the remote engine when it is not configured to handle everything.
A root document with versions reprocesses its latest version, which is the
file whose content, archive and thumbnail are shown for it.
""" """
for document_id in doc_ids: latest_version = versions_newest_first(
Document.objects.filter(root_document=OuterRef("pk")),
).values("id")[:1]
source_ids = (
Document.objects.filter(id__in=doc_ids)
.annotate(
source_id=Case(
When(root_document__isnull=False, then=F("id")),
default=Coalesce(Subquery(latest_version), F("id")),
),
)
.order_by()
.values_list("source_id", flat=True)
.distinct()
)
for source_id in source_ids:
update_document_content_maybe_archive_file.apply_async( update_document_content_maybe_archive_file.apply_async(
kwargs={"document_id": document_id, "remote_ocr": remote_ocr}, kwargs={"document_id": source_id, "remote_ocr": remote_ocr},
headers={"trigger_source": PaperlessTask.TriggerSource.MANUAL}, headers={"trigger_source": PaperlessTask.TriggerSource.MANUAL},
) )
+38 -32
View File
@@ -31,6 +31,7 @@ from documents.search._query import parse_user_query
from documents.search._schema import _write_sentinels from documents.search._schema import _write_sentinels
from documents.search._schema import build_schema from documents.search._schema import build_schema
from documents.search._schema import open_or_rebuild_index from documents.search._schema import open_or_rebuild_index
from documents.search._schema import rebuild_in_progress
from documents.search._schema import wipe_index from documents.search._schema import wipe_index
from documents.search._tokenizer import ascii_fold from documents.search._tokenizer import ascii_fold
from documents.search._tokenizer import autocomplete_tokens from documents.search._tokenizer import autocomplete_tokens
@@ -1110,39 +1111,44 @@ class TantivyBackend:
flushing a segment, deferring merge work; they do not avoid it. flushing a segment, deferring merge work; they do not avoid it.
""" """
wipe_index(self._path) wipe_index(self._path)
new_index = tantivy.Index(build_schema(), path=str(self._path)) # The marker covers the window where the empty index is already stamped
_write_sentinels(self._path) # as current but not yet populated, so an interrupted rebuild is retried.
register_tokenizers(new_index, settings.SEARCH_LANGUAGE) with rebuild_in_progress(self._path):
new_index = tantivy.Index(build_schema(), path=str(self._path))
_write_sentinels(self._path)
register_tokenizers(new_index, settings.SEARCH_LANGUAGE)
# Point instance at the new index so _build_tantivy_doc uses it # Point instance at the new index so _build_tantivy_doc uses it
old_index, old_schema = self._raw_index, self._raw_schema old_index, old_schema = self._raw_index, self._raw_schema
self._raw_index = new_index self._raw_index = new_index
self._raw_schema = new_index.schema self._raw_schema = new_index.schema
# Stream documents one-by-one (so the progress bar advances per # Stream documents one-by-one (so the progress bar advances per
# document) while fetching viewer permissions one SQL query per chunk. # document) while fetching viewer permissions one SQL query per
# The stream is Sized, so iter_wrapper can still discover the total. # chunk. The stream is Sized, so iter_wrapper can still discover
documents_stream = _DocumentViewerStream(documents, chunk_size=1000) # the total.
try: documents_stream = _DocumentViewerStream(documents, chunk_size=1000)
writer = new_index.writer(heap_size=writer_heap_bytes) try:
for document, (viewer_ids, viewer_group_ids) in iter_wrapper( writer = new_index.writer(heap_size=writer_heap_bytes)
documents_stream, for document, (viewer_ids, viewer_group_ids) in iter_wrapper(
): documents_stream,
doc = self._build_tantivy_doc( ):
document, doc = self._build_tantivy_doc(
viewer_ids=viewer_ids, document,
viewer_group_ids=viewer_group_ids, viewer_ids=viewer_ids,
) viewer_group_ids=viewer_group_ids,
writer.add_document(doc) )
writer.commit() writer.add_document(doc)
# Wait for background merge threads to finish so all segments are writer.commit()
# fully merged and persisted before the index is considered rebuilt. # Wait for background merge threads to finish so all segments
writer.wait_merging_threads() # are fully merged and persisted before the index is considered
new_index.reload() # rebuilt.
except BaseException: # pragma: no cover writer.wait_merging_threads()
# Restore old index on failure so the backend remains usable new_index.reload()
self._raw_index = old_index except BaseException: # pragma: no cover
self._raw_schema = old_schema # Restore old index on failure so the backend remains usable
raise self._raw_index = old_index
self._raw_schema = old_schema
raise
def chunked(iterable, size): def chunked(iterable, size):
+45 -4
View File
@@ -4,6 +4,7 @@ import hashlib
import json import json
import logging import logging
import shutil import shutil
from contextlib import contextmanager
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
from typing import Final from typing import Final
from typing import NamedTuple from typing import NamedTuple
@@ -16,6 +17,7 @@ from whoosh_compat import FieldKind
from documents.search._fields import PUBLIC_FIELDS from documents.search._fields import PUBLIC_FIELDS
if TYPE_CHECKING: if TYPE_CHECKING:
from collections.abc import Iterator
from pathlib import Path from pathlib import Path
logger = logging.getLogger("paperless.search") logger = logging.getLogger("paperless.search")
@@ -28,6 +30,11 @@ logger = logging.getLogger("paperless.search")
# v3 - barcodes JSON field for stored barcode contents # v3 - barcodes JSON field for stored barcode contents
SCHEMA_VERSION: Final[int] = 3 SCHEMA_VERSION: Final[int] = 3
# Present in the index directory from the moment a full rebuild starts until it
# finishes. If a rebuild is interrupted it is left behind, so the half-built
# index is not mistaken for a complete one.
REBUILD_MARKER: Final[str] = ".rebuilding"
class FieldDescriptor(NamedTuple): class FieldDescriptor(NamedTuple):
"""One tantivy field, in declaration order. """One tantivy field, in declaration order.
@@ -255,9 +262,9 @@ def needs_rebuild(index_dir: Path) -> bool:
""" """
Check if the search index needs rebuilding. Check if the search index needs rebuilding.
Reads .index_settings.json to compare the stored schema version, search True if a previous full rebuild never finished (the rebuild marker is still
language and schema fingerprint against the current configuration. Returns present), or if the index's stamped settings no longer match the current
True if the file is missing, unparsable, or any value mismatches. configuration. See _settings_mismatch().
Args: Args:
index_dir: Path to the search index directory index_dir: Path to the search index directory
@@ -265,6 +272,40 @@ def needs_rebuild(index_dir: Path) -> bool:
Returns: Returns:
True if the index needs rebuilding, False if it's up to date True if the index needs rebuilding, False if it's up to date
""" """
if (index_dir / REBUILD_MARKER).exists():
logger.warning("Previous search index rebuild did not finish - rebuilding.")
return True
return _settings_mismatch(index_dir)
@contextmanager
def rebuild_in_progress(index_dir: Path) -> Iterator[None]:
"""
Flag the index as incomplete for the duration of a full rebuild.
The marker is cleared only if the block exits cleanly. There is deliberately
no try/finally: an exception must leave the marker behind so the next
needs_rebuild() check retries the rebuild.
"""
marker = index_dir / REBUILD_MARKER
marker.touch()
yield
marker.unlink(missing_ok=True)
def _settings_mismatch(index_dir: Path) -> bool:
"""
Check the stamped settings against the current configuration.
Reads .index_settings.json to compare the stored schema version, search
language and schema fingerprint. Returns True if the file is missing,
unparsable, or any value mismatches.
This deliberately ignores the rebuild marker: open_or_rebuild_index() uses it
so that a process opening the index while another process is mid-rebuild
(or after one died) does not wipe the partial index out from under it.
Repopulating is the job of ``document_index reindex``.
"""
settings_file = index_dir / ".index_settings.json" settings_file = index_dir / ".index_settings.json"
if not settings_file.exists(): if not settings_file.exists():
return True return True
@@ -333,7 +374,7 @@ def open_or_rebuild_index(index_dir: Path | None = None) -> tantivy.Index:
index_dir = cast("Path", settings.INDEX_DIR) index_dir = cast("Path", settings.INDEX_DIR)
if not index_dir.exists(): if not index_dir.exists():
return tantivy.Index(build_schema()) return tantivy.Index(build_schema())
if needs_rebuild(index_dir): if _settings_mismatch(index_dir):
wipe_index(index_dir) wipe_index(index_dir)
idx = tantivy.Index(build_schema(), path=str(index_dir)) idx = tantivy.Index(build_schema(), path=str(index_dir))
_write_sentinels(index_dir) _write_sentinels(index_dir)
+8 -3
View File
@@ -490,18 +490,23 @@ def update_document_content_maybe_archive_file(
shutil.move(thumbnail, document.thumbnail_path) shutil.move(thumbnail, document.thumbnail_path)
document.refresh_from_db() document.refresh_from_db()
root_document = (
document.root_document if document.root_document_id else document
)
logger.info( logger.info(
f"Updating index for document {document_id} ({document.archive_checksum})", f"Updating index for document {root_document.pk} ({document.archive_checksum})",
) )
from documents.search import get_backend from documents.search import get_backend
get_backend().add_or_update(document) get_backend().add_or_update(root_document)
ai_config = AIConfig() ai_config = AIConfig()
if ai_config.llm_index_enabled: if ai_config.llm_index_enabled:
llm_index_add_or_update_document(document) llm_index_add_or_update_document(root_document)
clear_document_caches(document.pk) clear_document_caches(document.pk)
if root_document.pk != document.pk:
clear_document_caches(root_document.pk)
except Exception: except Exception:
logger.exception( logger.exception(
@@ -16,7 +16,10 @@ from documents.search._backend import TantivyBackend
from documents.search._backend import WriteBatch from documents.search._backend import WriteBatch
from documents.search._backend import get_backend from documents.search._backend import get_backend
from documents.search._backend import reset_backend from documents.search._backend import reset_backend
from documents.search._schema import REBUILD_MARKER
from documents.search._schema import needs_rebuild
from documents.signals.handlers import add_to_index from documents.signals.handlers import add_to_index
from paperless_testing.dirs import PaperlessDirs
from paperless_testing.factories import CorrespondentFactory from paperless_testing.factories import CorrespondentFactory
from paperless_testing.factories import DocumentFactory from paperless_testing.factories import DocumentFactory
from paperless_testing.factories import DocumentTypeFactory from paperless_testing.factories import DocumentTypeFactory
@@ -823,6 +826,53 @@ class TestRebuild:
backend.rebuild(Document.objects.all(), iter_wrapper=wrapper) backend.rebuild(Document.objects.all(), iter_wrapper=wrapper)
assert 30 in seen assert 30 in seen
def test_successful_rebuild_leaves_index_up_to_date(
self,
backend: TantivyBackend,
paperless_dirs: PaperlessDirs,
) -> None:
"""
GIVEN:
- A backend and one document
WHEN:
- rebuild() completes
THEN:
- needs_rebuild() is False and no rebuild marker remains
"""
DocumentFactory.create()
backend.rebuild(Document.objects.all())
assert needs_rebuild(paperless_dirs.index_dir) is False
assert not (paperless_dirs.index_dir / REBUILD_MARKER).exists()
def test_interrupted_rebuild_is_retried(
self,
backend: TantivyBackend,
paperless_dirs: PaperlessDirs,
) -> None:
"""
GIVEN:
- A rebuild that dies while indexing documents (e.g. the database
connection is lost)
WHEN:
- needs_rebuild() is checked afterwards
THEN:
- It is True, even though the empty index was already stamped with
current settings, so the next start rebuilds instead of reporting
the index as up to date
"""
DocumentFactory.create()
def die(pairs):
raise RuntimeError("terminating connection due to administrator command")
yield # pragma: no cover
with pytest.raises(RuntimeError):
backend.rebuild(Document.objects.all(), iter_wrapper=die)
assert needs_rebuild(paperless_dirs.index_dir) is True
def test_includes_group_granted_viewers(self, backend: TantivyBackend) -> None: def test_includes_group_granted_viewers(self, backend: TantivyBackend) -> None:
"""Rebuild must index viewer ids for group-only grants, not just direct ones. """Rebuild must index viewer ids for group-only grants, not just direct ones.
+91 -27
View File
@@ -4,6 +4,7 @@ from pathlib import Path
from unittest import mock from unittest import mock
import pikepdf import pikepdf
import pytest
from django.contrib.auth.models import Group from django.contrib.auth.models import Group
from django.contrib.auth.models import Permission from django.contrib.auth.models import Permission
from django.contrib.auth.models import User from django.contrib.auth.models import User
@@ -12,6 +13,7 @@ from django.test import TestCase
from django.test.utils import CaptureQueriesContext from django.test.utils import CaptureQueriesContext
from guardian.shortcuts import get_groups_with_perms from guardian.shortcuts import get_groups_with_perms
from guardian.shortcuts import get_users_with_perms from guardian.shortcuts import get_users_with_perms
from pytest_mock import MockerFixture
from documents import bulk_edit from documents import bulk_edit
from documents.models import Correspondent from documents.models import Correspondent
@@ -23,6 +25,7 @@ from documents.models import StoragePath
from documents.models import Tag from documents.models import Tag
from documents.permissions import set_permissions_for_objects from documents.permissions import set_permissions_for_objects
from paperless_testing.dirs import DirectoriesMixin from paperless_testing.dirs import DirectoriesMixin
from paperless_testing.factories import DocumentFactory
from paperless_testing.permissions import grant_object from paperless_testing.permissions import grant_object
@@ -1970,18 +1973,22 @@ class TestPDFActions(DirectoriesMixin, TestCase):
self.assertIn("Error removing password from document", cm.output[0]) self.assertIn("Error removing password from document", cm.output[0])
class TestBulkEditReprocess(DirectoriesMixin, TestCase): @pytest.mark.django_db
def setUp(self) -> None: class TestBulkEditReprocess:
super().setUp() @pytest.fixture
def mock_task(self, mocker: MockerFixture) -> mock.MagicMock:
self.doc = Document.objects.create( return mocker.patch(
title="test", "documents.bulk_edit.update_document_content_maybe_archive_file",
checksum="A",
mime_type="application/pdf",
) )
@mock.patch("documents.bulk_edit.update_document_content_maybe_archive_file") @staticmethod
def test_reprocess_defaults_to_local(self, mock_task: mock.Mock) -> None: def _queued_ids(mock_task: mock.MagicMock) -> list[int]:
return [
call.kwargs["kwargs"]["document_id"]
for call in mock_task.apply_async.call_args_list
]
def test_reprocess_defaults_to_local(self, mock_task: mock.MagicMock) -> None:
""" """
GIVEN: GIVEN:
- A reprocess request that says nothing about remote OCR - A reprocess request that says nothing about remote OCR
@@ -1990,18 +1997,17 @@ class TestBulkEditReprocess(DirectoriesMixin, TestCase):
THEN: THEN:
- The task is queued without asking for the remote engine - The task is queued without asking for the remote engine
""" """
result = bulk_edit.reprocess([self.doc.id]) doc = DocumentFactory()
assert bulk_edit.reprocess([doc.id]) == "OK"
self.assertEqual(result, "OK")
mock_task.apply_async.assert_called_once() mock_task.apply_async.assert_called_once()
_, kwargs = mock_task.apply_async.call_args assert mock_task.apply_async.call_args.kwargs["kwargs"] == {
self.assertEqual( "document_id": doc.id,
kwargs["kwargs"], "remote_ocr": False,
{"document_id": self.doc.id, "remote_ocr": False}, }
)
@mock.patch("documents.bulk_edit.update_document_content_maybe_archive_file") def test_reprocess_passes_remote_ocr(self, mock_task: mock.MagicMock) -> None:
def test_reprocess_passes_remote_ocr(self, mock_task: mock.Mock) -> None:
""" """
GIVEN: GIVEN:
- A reprocess request that explicitly asks for remote OCR - A reprocess request that explicitly asks for remote OCR
@@ -2010,14 +2016,72 @@ class TestBulkEditReprocess(DirectoriesMixin, TestCase):
THEN: THEN:
- The request is forwarded to the task for every document - The request is forwarded to the task for every document
""" """
other = Document.objects.create( docs = DocumentFactory.create_batch(2)
title="test2",
checksum="B", bulk_edit.reprocess([doc.id for doc in docs], remote_ocr=True)
mime_type="application/pdf",
assert mock_task.apply_async.call_count == 2
for call in mock_task.apply_async.call_args_list:
assert call.kwargs["kwargs"]["remote_ocr"]
def test_reprocess_root_uses_latest_version(
self,
mock_task: mock.MagicMock,
) -> None:
"""
GIVEN:
- A root document with two versions
WHEN:
- reprocess is called with the root document
THEN:
- The latest version is reprocessed, not the root's original file
"""
root = DocumentFactory()
DocumentFactory(root_document=root, version_index=1)
latest = DocumentFactory(root_document=root, version_index=2)
bulk_edit.reprocess([root.id])
assert self._queued_ids(mock_task) == [latest.id]
def test_reprocess_explicit_version(self, mock_task: mock.MagicMock) -> None:
"""
GIVEN:
- A root document with two versions
WHEN:
- reprocess is called with the older version
THEN:
- That version is reprocessed
"""
root = DocumentFactory()
older = DocumentFactory(root_document=root, version_index=1)
DocumentFactory(root_document=root, version_index=2)
bulk_edit.reprocess([older.id])
assert self._queued_ids(mock_task) == [older.id]
def test_reprocess_root_and_latest_version_dispatches_once(
self,
mock_task: mock.MagicMock,
) -> None:
"""
GIVEN:
- A root document with two versions, the latest created on a
different date than the root
WHEN:
- reprocess is called with both the root and its latest version
THEN:
- The latest version is reprocessed only once
"""
root = DocumentFactory(created=date(2024, 1, 1))
DocumentFactory(root_document=root, version_index=1)
latest = DocumentFactory(
root_document=root,
version_index=2,
created=date(2025, 1, 1),
) )
bulk_edit.reprocess([self.doc.id, other.id], remote_ocr=True) bulk_edit.reprocess([root.id, latest.id])
self.assertEqual(mock_task.apply_async.call_count, 2) assert self._queued_ids(mock_task) == [latest.id]
for call in mock_task.apply_async.call_args_list:
self.assertTrue(call.kwargs["kwargs"]["remote_ocr"])
+77
View File
@@ -20,6 +20,7 @@ from documents.sanity_checker import SanityCheckMessages
from documents.tests.helpers import dummy_preprocess from documents.tests.helpers import dummy_preprocess
from paperless_testing.assertions import FileSystemAssertsMixin from paperless_testing.assertions import FileSystemAssertsMixin
from paperless_testing.dirs import DirectoriesMixin from paperless_testing.dirs import DirectoriesMixin
from paperless_testing.factories import DocumentFactory
@pytest.mark.django_db @pytest.mark.django_db
@@ -287,6 +288,82 @@ class TestUpdateContent(DirectoriesMixin, TestCase):
tasks.update_document_content_maybe_archive_file(doc.pk) tasks.update_document_content_maybe_archive_file(doc.pk)
self.assertNotEqual(Document.objects.get(pk=doc.pk).content, "test") self.assertNotEqual(Document.objects.get(pk=doc.pk).content, "test")
def _create_root_with_version(self) -> tuple[Document, Document]:
sample1 = self.dirs.scratch_dir / "sample.pdf"
shutil.copy(
Path(__file__).parent
/ "samples"
/ "documents"
/ "originals"
/ "0000001.pdf",
sample1,
)
root = DocumentFactory(content="root content", mime_type="application/pdf")
version = DocumentFactory(
content="my document",
filename=sample1,
mime_type="application/pdf",
root_document=root,
version_index=1,
)
return root, version
@mock.patch("documents.tasks.clear_document_caches")
@mock.patch("documents.search.get_backend")
def test_update_content_version_indexes_root(
self,
mock_get_backend: mock.Mock,
mock_clear_caches: mock.Mock,
) -> None:
"""
GIVEN:
- A root document with a version
WHEN:
- Update content task is called for the version
THEN:
- The version's content is updated
- The root document is indexed rather than the version
- Caches are cleared for both
"""
root, version = self._create_root_with_version()
tasks.update_document_content_maybe_archive_file(version.pk)
self.assertNotEqual(
Document.objects.get(pk=version.pk).content,
"my document",
)
self.assertEqual(Document.objects.get(pk=root.pk).content, "root content")
indexed = mock_get_backend.return_value.add_or_update.call_args.args[0]
self.assertEqual(indexed.pk, root.pk)
mock_clear_caches.assert_has_calls(
[mock.call(version.pk), mock.call(root.pk)],
)
@override_settings(AI_ENABLED=True, LLM_EMBEDDING_BACKEND="huggingface")
@mock.patch("documents.tasks.llm_index_add_or_update_document")
@mock.patch("documents.search.get_backend")
def test_update_content_version_updates_llm_index_for_root(
self,
mock_get_backend: mock.Mock,
mock_llm_index: mock.Mock,
) -> None:
"""
GIVEN:
- A root document with a version
- The LLM index is enabled
WHEN:
- Update content task is called for the version
THEN:
- The LLM index is updated for the root document, not the version
"""
root, version = self._create_root_with_version()
tasks.update_document_content_maybe_archive_file(version.pk)
mock_llm_index.assert_called_once()
self.assertEqual(mock_llm_index.call_args.args[0].pk, root.pk)
class TestUpdateContentRemoteOCR(DirectoriesMixin, TestCase): class TestUpdateContentRemoteOCR(DirectoriesMixin, TestCase):
""" """