diff --git a/src/documents/search/_backend.py b/src/documents/search/_backend.py index 996d332e8..be9376ea4 100644 --- a/src/documents/search/_backend.py +++ b/src/documents/search/_backend.py @@ -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): diff --git a/src/documents/search/_schema.py b/src/documents/search/_schema.py index 2539d12fe..62130c246 100644 --- a/src/documents/search/_schema.py +++ b/src/documents/search/_schema.py @@ -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) diff --git a/src/documents/tests/search/test_backend.py b/src/documents/tests/search/test_backend.py index d215e3ca7..223907e19 100644 --- a/src/documents/tests/search/test_backend.py +++ b/src/documents/tests/search/test_backend.py @@ -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.