Compare commits

...
Author SHA1 Message Date
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 270 additions and 41 deletions

No files matched your search

+8 -2
View File
@@ -408,10 +408,16 @@ 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
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:
for doc in Document.objects.select_related("root_document").filter(
id__in=doc_ids,
):
pair = _resolve_root_and_source_doc(doc)
update_document_content_maybe_archive_file.apply_async(
kwargs={"document_id": document_id, "remote_ocr": remote_ocr},
kwargs={"document_id": pair.source_doc.id, "remote_ocr": remote_ocr},
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 build_schema
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._tokenizer import ascii_fold
from documents.search._tokenizer import autocomplete_tokens
@@ -1110,39 +1111,44 @@ class TantivyBackend:
flushing a segment, deferring merge work; they do not avoid it.
"""
wipe_index(self._path)
new_index = tantivy.Index(build_schema(), path=str(self._path))
_write_sentinels(self._path)
register_tokenizers(new_index, settings.SEARCH_LANGUAGE)
# The marker covers the window where the empty index is already stamped
# as current but not yet populated, so an interrupted rebuild is retried.
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
old_index, old_schema = self._raw_index, self._raw_schema
self._raw_index = new_index
self._raw_schema = new_index.schema
# Stream documents one-by-one (so the progress bar advances per
# document) while fetching viewer permissions one SQL query per chunk.
# The stream is Sized, so iter_wrapper can still discover the total.
documents_stream = _DocumentViewerStream(documents, chunk_size=1000)
try:
writer = new_index.writer(heap_size=writer_heap_bytes)
for document, (viewer_ids, viewer_group_ids) in iter_wrapper(
documents_stream,
):
doc = self._build_tantivy_doc(
document,
viewer_ids=viewer_ids,
viewer_group_ids=viewer_group_ids,
)
writer.add_document(doc)
writer.commit()
# Wait for background merge threads to finish so all segments are
# fully merged and persisted before the index is considered rebuilt.
writer.wait_merging_threads()
new_index.reload()
except BaseException: # pragma: no cover
# Restore old index on failure so the backend remains usable
self._raw_index = old_index
self._raw_schema = old_schema
raise
# Point instance at the new index so _build_tantivy_doc uses it
old_index, old_schema = self._raw_index, self._raw_schema
self._raw_index = new_index
self._raw_schema = new_index.schema
# Stream documents one-by-one (so the progress bar advances per
# document) while fetching viewer permissions one SQL query per
# chunk. The stream is Sized, so iter_wrapper can still discover
# the total.
documents_stream = _DocumentViewerStream(documents, chunk_size=1000)
try:
writer = new_index.writer(heap_size=writer_heap_bytes)
for document, (viewer_ids, viewer_group_ids) in iter_wrapper(
documents_stream,
):
doc = self._build_tantivy_doc(
document,
viewer_ids=viewer_ids,
viewer_group_ids=viewer_group_ids,
)
writer.add_document(doc)
writer.commit()
# Wait for background merge threads to finish so all segments
# are fully merged and persisted before the index is considered
# rebuilt.
writer.wait_merging_threads()
new_index.reload()
except BaseException: # pragma: no cover
# Restore old index on failure so the backend remains usable
self._raw_index = old_index
self._raw_schema = old_schema
raise
def chunked(iterable, size):
+45 -4
View File
@@ -4,6 +4,7 @@ import hashlib
import json
import logging
import shutil
from contextlib import contextmanager
from typing import TYPE_CHECKING
from typing import Final
from typing import NamedTuple
@@ -16,6 +17,7 @@ from whoosh_compat import FieldKind
from documents.search._fields import PUBLIC_FIELDS
if TYPE_CHECKING:
from collections.abc import Iterator
from pathlib import Path
logger = logging.getLogger("paperless.search")
@@ -28,6 +30,11 @@ logger = logging.getLogger("paperless.search")
# v3 - barcodes JSON field for stored barcode contents
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):
"""One tantivy field, in declaration order.
@@ -255,9 +262,9 @@ def needs_rebuild(index_dir: Path) -> bool:
"""
Check if the search index needs rebuilding.
Reads .index_settings.json to compare the stored schema version, search
language and schema fingerprint against the current configuration. Returns
True if the file is missing, unparsable, or any value mismatches.
True if a previous full rebuild never finished (the rebuild marker is still
present), or if the index's stamped settings no longer match the current
configuration. See _settings_mismatch().
Args:
index_dir: Path to the search index directory
@@ -265,6 +272,40 @@ def needs_rebuild(index_dir: Path) -> bool:
Returns:
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"
if not settings_file.exists():
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)
if not index_dir.exists():
return tantivy.Index(build_schema())
if needs_rebuild(index_dir):
if _settings_mismatch(index_dir):
wipe_index(index_dir)
idx = tantivy.Index(build_schema(), path=str(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)
document.refresh_from_db()
root_document = (
document.root_document if document.root_document_id else document
)
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
get_backend().add_or_update(document)
get_backend().add_or_update(root_document)
ai_config = AIConfig()
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)
if root_document.pk != document.pk:
clear_document_caches(root_document.pk)
except Exception:
logger.exception(
@@ -16,7 +16,10 @@ from documents.search._backend import TantivyBackend
from documents.search._backend import WriteBatch
from documents.search._backend import get_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 paperless_testing.dirs import PaperlessDirs
from paperless_testing.factories import CorrespondentFactory
from paperless_testing.factories import DocumentFactory
from paperless_testing.factories import DocumentTypeFactory
@@ -823,6 +826,53 @@ class TestRebuild:
backend.rebuild(Document.objects.all(), iter_wrapper=wrapper)
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:
"""Rebuild must index viewer ids for group-only grants, not just direct ones.
+66
View File
@@ -2021,3 +2021,69 @@ class TestBulkEditReprocess(DirectoriesMixin, TestCase):
self.assertEqual(mock_task.apply_async.call_count, 2)
for call in mock_task.apply_async.call_args_list:
self.assertTrue(call.kwargs["kwargs"]["remote_ocr"])
@mock.patch("documents.bulk_edit.update_document_content_maybe_archive_file")
def test_reprocess_root_uses_latest_version(self, mock_task: mock.Mock) -> 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
"""
Document.objects.create(
title="test",
checksum="A-v1",
mime_type="application/pdf",
root_document=self.doc,
version_index=1,
)
latest = Document.objects.create(
title="test",
checksum="A-v2",
mime_type="application/pdf",
root_document=self.doc,
version_index=2,
)
bulk_edit.reprocess([self.doc.id])
mock_task.apply_async.assert_called_once()
self.assertEqual(
mock_task.apply_async.call_args.kwargs["kwargs"]["document_id"],
latest.id,
)
@mock.patch("documents.bulk_edit.update_document_content_maybe_archive_file")
def test_reprocess_explicit_version(self, mock_task: mock.Mock) -> None:
"""
GIVEN:
- A root document with two versions
WHEN:
- reprocess is called with the older version
THEN:
- That version is reprocessed
"""
older = Document.objects.create(
title="test",
checksum="A-v1",
mime_type="application/pdf",
root_document=self.doc,
version_index=1,
)
Document.objects.create(
title="test",
checksum="A-v2",
mime_type="application/pdf",
root_document=self.doc,
version_index=2,
)
bulk_edit.reprocess([older.id])
mock_task.apply_async.assert_called_once()
self.assertEqual(
mock_task.apply_async.call_args.kwargs["kwargs"]["document_id"],
older.id,
)
+55
View File
@@ -287,6 +287,61 @@ class TestUpdateContent(DirectoriesMixin, TestCase):
tasks.update_document_content_maybe_archive_file(doc.pk)
self.assertNotEqual(Document.objects.get(pk=doc.pk).content, "test")
@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
"""
sample1 = self.dirs.scratch_dir / "sample.pdf"
shutil.copy(
Path(__file__).parent
/ "samples"
/ "documents"
/ "originals"
/ "0000001.pdf",
sample1,
)
root = Document.objects.create(
title="test",
content="root content",
checksum="root",
mime_type="application/pdf",
)
version = Document.objects.create(
title="test",
content="my document",
checksum="wow",
filename=sample1,
mime_type="application/pdf",
root_document=root,
version_index=1,
)
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)],
)
class TestUpdateContentRemoteOCR(DirectoriesMixin, TestCase):
"""