Coverage for documents/search/_backend.py: 37%
478 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 random
5import re
6import threading
7import time
8from datetime import UTC
9from datetime import datetime
10from enum import StrEnum
11from itertools import islice
12from typing import TYPE_CHECKING
13from typing import Final
14from typing import NamedTuple
15from typing import Self
16from typing import TypedDict
17from typing import TypeVar
18from typing import cast
20import filelock
21import tantivy
22from django.conf import settings
23from django.utils.timezone import get_current_timezone
25from documents.search._query import extract_cjk_text
26from documents.search._query import normalize_search_text
27from documents.search._query import parse_simple_text_highlight_query
28from documents.search._query import parse_simple_text_query
29from documents.search._query import parse_simple_title_query
30from documents.search._query import parse_user_query
31from documents.search._schema import _write_sentinels
32from documents.search._schema import build_schema
33from documents.search._schema import open_or_rebuild_index
34from documents.search._schema import wipe_index
35from documents.search._tokenizer import ascii_fold
36from documents.search._tokenizer import autocomplete_tokens
37from documents.search._tokenizer import register_tokenizers
38from documents.utils import IterWrapper
39from documents.utils import QuerySetStream
40from documents.utils import identity
42if TYPE_CHECKING: 42 ↛ 43line 42 didn't jump to line 43 because the condition on line 42 was never true
43 from collections.abc import Iterable
44 from collections.abc import Iterator
45 from collections.abc import Sequence
46 from pathlib import Path
48 from django.contrib.auth.models import AbstractUser
49 from django.contrib.auth.models import Group
50 from django.contrib.auth.models import User
51 from django.db.models import QuerySet
53 from documents.models import Document
55logger = logging.getLogger("paperless.search")
57_LOCK_TIMEOUT_SECONDS: Final[float] = 10.0 # per-attempt acquire timeout
58_LOCK_RETRY_ATTEMPTS: Final[int] = 4 # total attempts (1 initial + 3 retries)
59_LOCK_BACKOFF_BASE: Final[float] = 1.0 # seconds
60_LOCK_BACKOFF_CAP: Final[float] = 10.0 # seconds
62T = TypeVar("T")
65class ViewerGrant(NamedTuple):
66 """Direct user and group view grants for a single document.
68 Named fields (rather than a bare 2-tuple) so ``viewer_ids`` and
69 ``viewer_group_ids`` can't be silently transposed at a call site — both
70 are ``list[int]``, so a positional swap would type-check cleanly.
71 """
73 viewer_ids: list[int]
74 viewer_group_ids: list[int]
77class SearchMode(StrEnum):
78 QUERY = "query"
79 TEXT = "text"
80 TITLE = "title"
83def _extract_autocomplete_words(text_sources: list[str]) -> set[str]:
84 """Extract and normalize words for autocomplete.
86 Tokenizes with Tantivy's simple analyzer (simple -> lowercase -> ascii_fold)
87 so the extracted words match how document content is indexed, and runs the
88 whole pass in Rust. This replaces a Python regex scan plus per-token folding,
89 which dominated full-reindex CPU time. Tokenizing natively also removes the
90 ReDoS exposure of running a regex over untrusted content.
91 """
92 words = set()
93 for text in text_sources:
94 if text:
95 words.update(autocomplete_tokens(text))
96 return words
99class SearchHit(TypedDict):
100 """Type definition for search result hits."""
102 id: int
103 score: float
104 rank: int
105 highlights: dict[str, str]
108class TantivyRelevanceList:
109 """
110 DRF-compatible list wrapper for Tantivy search results.
112 Holds a lightweight ordered list of IDs (for pagination count and
113 ``selection_data``) together with a small page of rich ``SearchHit``
114 dicts (for serialization). DRF's ``PageNumberPagination`` calls
115 ``__len__`` to compute the total page count and ``__getitem__`` to
116 slice the displayed page.
118 Args:
119 ordered_ids: All matching document IDs in display order.
120 page_hits: Rich SearchHit dicts for the requested DRF page only.
121 page_offset: Index into *ordered_ids* where *page_hits* starts.
122 """
124 def __init__(
125 self,
126 ordered_ids: list[int],
127 page_hits: list[SearchHit],
128 page_offset: int = 0,
129 ) -> None:
130 self._ordered_ids = ordered_ids
131 self._page_hits = page_hits
132 self._page_offset = page_offset
134 def __len__(self) -> int:
135 return len(self._ordered_ids)
137 def __getitem__(self, key: int | slice) -> SearchHit | list[SearchHit]:
138 if isinstance(key, int): 138 ↛ 139line 138 didn't jump to line 139 because the condition on line 138 was never true
139 idx = key if key >= 0 else len(self._ordered_ids) + key
140 if self._page_offset <= idx < self._page_offset + len(self._page_hits):
141 return self._page_hits[idx - self._page_offset]
142 return SearchHit(
143 id=self._ordered_ids[key],
144 score=0.0,
145 rank=idx + 1,
146 highlights={},
147 )
148 start = key.start or 0
149 stop = key.stop or len(self._ordered_ids)
150 # DRF slices to extract the current page. If the slice aligns
151 # with our pre-fetched page_hits, return them directly.
152 # We only check start — DRF always slices with stop=start+page_size,
153 # which exceeds page_hits length on the last page.
154 if start == self._page_offset: 154 ↛ 157line 154 didn't jump to line 157 because the condition on line 154 was always true
155 return self._page_hits[: stop - start]
156 # Fallback: return stub dicts (no highlights).
157 return [
158 SearchHit(id=doc_id, score=0.0, rank=start + i + 1, highlights={})
159 for i, doc_id in enumerate(self._ordered_ids[key])
160 ]
162 def get_all_ids(self) -> list[int]:
163 """Return all matching document IDs in display order."""
164 return self._ordered_ids
167class SearchIndexLockError(Exception):
168 """Raised when the search index file lock cannot be acquired within the timeout."""
171class WriteBatch:
172 """
173 Context manager for bulk index operations with file locking.
175 Provides transactional batch updates to the search index with proper
176 concurrency control via file locking. All operations within the batch
177 are committed atomically or rolled back on exception.
179 Usage:
180 with backend.batch_update() as batch:
181 batch.add_or_update(document)
182 batch.remove(doc_id)
183 """
185 def __init__(self, backend: TantivyBackend, lock_timeout: float):
186 self._backend = backend
187 self._lock_timeout = lock_timeout
188 self._raw_writer: tantivy.IndexWriter | None = None
189 self._lock = None
191 @property
192 def _writer(self) -> tantivy.IndexWriter:
193 assert self._raw_writer is not None, (
194 "WriteBatch not entered; use as context manager"
195 )
196 return self._raw_writer
198 def __enter__(self) -> Self:
199 lock_path = self._backend._path / ".tantivy.lock"
200 self._lock = filelock.FileLock(str(lock_path))
201 for attempt in range(_LOCK_RETRY_ATTEMPTS): 201 ↛ 236line 201 didn't jump to line 236 because the loop on line 201 didn't complete
202 try:
203 self._lock.acquire(timeout=self._lock_timeout)
204 break
205 except filelock.Timeout:
206 if attempt == _LOCK_RETRY_ATTEMPTS - 1:
207 raise SearchIndexLockError(
208 f"Could not acquire index lock after {_LOCK_RETRY_ATTEMPTS} "
209 f"attempts (timeout={self._lock_timeout}s each)",
210 )
211 sleep_s = random.uniform(
212 0,
213 min(_LOCK_BACKOFF_CAP, _LOCK_BACKOFF_BASE * (2**attempt)),
214 )
215 logger.debug(
216 "Index lock contention; retrying in %.2fs (attempt %d/%d)",
217 sleep_s,
218 attempt + 1,
219 _LOCK_RETRY_ATTEMPTS,
220 )
221 time.sleep(sleep_s)
223 # Open a fresh Index (and thus a fresh Tantivy ManagedDirectory)
224 # for the write, rather than reusing the process-local cached
225 # index. ManagedDirectory loads its GC bookkeeping (.managed.json)
226 # once, at construction, and never re-reads it; paperless runs
227 # several long-lived processes (Granian workers, Celery workers)
228 # that take turns writing under the file lock above. A cached,
229 # long-lived writer index would carry a stale managed-files view
230 # and, on commit, overwrite .managed.json with that stale view -
231 # permanently losing track of segment files other processes
232 # registered in the meantime, so they can never be garbage
233 # collected. Reopening fresh here always picks up the current
234 # on-disk state. The long-lived self._backend._index is used for
235 # reads only and is reloaded (not reopened) after commit below.
236 write_index = tantivy.Index(
237 build_schema(),
238 path=str(self._backend._path),
239 )
240 register_tokenizers(write_index, settings.SEARCH_LANGUAGE)
241 self._raw_writer = write_index.writer()
242 return self
244 def __exit__(self, exc_type, exc_val, exc_tb):
245 try:
246 if exc_type is None: 246 ↛ 258line 246 didn't jump to line 258 because the condition on line 246 was always true
247 self._writer.commit()
248 # Wait for background merge threads to finish before releasing
249 # the file lock so the next writer doesn't race against an
250 # in-progress merge on the same index files.
251 self._writer.wait_merging_threads()
252 self._backend._index.reload()
253 finally:
254 # Always release the writer (and Tantivy's internal writer lock),
255 # even if commit/merge/reload raised, so the next batch can acquire
256 # a writer instead of failing with LockBusy. An uncommitted writer
257 # is simply discarded.
258 if self._raw_writer is not None: 258 ↛ 261line 258 didn't jump to line 261 because the condition on line 258 was always true
259 del self._raw_writer
260 self._raw_writer = None
261 if self._lock is not None: 261 ↛ exitline 261 didn't return from function '__exit__' because the condition on line 261 was always true
262 self._lock.release()
264 def add_or_update(self, document: Document) -> None:
265 """
266 Add or update a document in the batch.
268 Implements upsert behavior by deleting any existing document with the same ID
269 and adding the new version. This ensures stale document data (e.g., after
270 permission changes) doesn't persist in the index.
272 Args:
273 document: Django Document instance to index
274 """
275 self.remove(document.pk)
276 doc = self._backend._build_tantivy_doc(document)
277 self._writer.add_document(doc)
279 def remove(self, doc_id: int) -> None:
280 """Remove a document from the batch by its primary key."""
281 self._writer.delete_documents_by_query(
282 tantivy.Query.term_query(self._backend._schema, "id", doc_id),
283 )
285 def add_or_update_ids(self, ids: Sequence[int]) -> None:
286 """
287 Add or update multiple documents in the batch by primary key.
289 Unlike calling ``add_or_update()`` once per document, this resolves
290 viewer permissions and effective (versioned) content in bulk against
291 the ids as a whole, instead of once per document -- see
292 ``_DocumentViewerStream`` and ``annotate_effective_content``. Use
293 this whenever more than one document is being written in the same
294 batch.
296 An id with no matching document (e.g. deleted between the caller
297 collecting ids and the batch running) is silently skipped, matching
298 ``add_or_update()``'s existing single-document deferred-task behavior
299 rather than erroring or leaving a stale index entry.
301 Args:
302 ids: Primary keys of Document instances to index
303 """
304 from documents.models import Document
305 from documents.versioning import annotate_effective_content
307 ids = list(ids)
308 if not ids:
309 return
311 queryset = annotate_effective_content(
312 Document.objects.filter(pk__in=ids)
313 .select_related("correspondent", "document_type", "storage_path", "owner")
314 .prefetch_related(
315 "tags",
316 "notes__user",
317 "custom_fields__field",
318 "barcodes",
319 "versions__barcodes",
320 ),
321 )
322 for document, grant in _DocumentViewerStream(queryset, chunk_size=1000):
323 self.remove(document.pk)
324 doc = self._backend._build_tantivy_doc(
325 document,
326 viewer_ids=grant.viewer_ids,
327 viewer_group_ids=grant.viewer_group_ids,
328 )
329 self._writer.add_document(doc)
332def build_permission_filter(
333 schema: tantivy.Schema,
334 user: AbstractUser,
335 viewer_group_ids: Iterable[int] = (),
336) -> tantivy.Query:
337 """
338 Build a query filter for user document permissions.
340 Creates a query that matches only documents visible to the specified user
341 according to paperless-ngx permission rules:
342 - Public documents (no owner) are visible to all users
343 - Private documents are visible to their owner
344 - Documents explicitly shared with the user are visible
345 - Documents shared with one of the user's current groups are visible
347 Args:
348 schema: Tantivy schema for field validation
349 user: User to check permissions for
350 viewer_group_ids: Current group memberships for the user
352 Returns:
353 Tantivy query that filters results to visible documents
354 """
355 owner_any = tantivy.Query.exists_query("owner_id")
356 no_owner = tantivy.Query.boolean_query(
357 [
358 (tantivy.Occur.Must, tantivy.Query.all_query()),
359 (tantivy.Occur.MustNot, owner_any),
360 ],
361 )
362 owned = tantivy.Query.term_query(schema, "owner_id", user.pk)
363 shared = tantivy.Query.term_query(schema, "viewer_id", user.pk)
364 group_shared = [
365 tantivy.Query.term_query(schema, "viewer_group_id", group_id)
366 for group_id in viewer_group_ids
367 ]
368 return tantivy.Query.disjunction_max_query(
369 [no_owner, owned, shared, *group_shared],
370 )
373class TantivyBackend:
374 """
375 Tantivy search backend with explicit lifecycle management.
377 Provides full-text search capabilities using the Tantivy search engine.
378 Keeps a persistent on-disk index. Handles document indexing, search queries,
379 autocompletion, and "more like this" functionality.
381 The backend manages its own connection lifecycle and can be reset when
382 the underlying index directory changes (e.g., during test isolation).
383 """
385 # Maps DRF ordering field names to Tantivy index field names.
386 SORT_FIELD_MAP: dict[str, str] = {
387 "title": "title_sort",
388 "correspondent__name": "correspondent_sort",
389 "document_type__name": "type_sort",
390 "created": "created",
391 "added": "added",
392 "modified": "modified",
393 "archive_serial_number": "asn",
394 "page_count": "page_count",
395 "num_notes": "num_notes",
396 }
398 # Fields where Tantivy's sort order matches the ORM's sort order.
399 # Text-based fields (title, correspondent__name, document_type__name)
400 # are excluded because Tantivy's tokenized fast fields produce different
401 # ordering than the ORM's collation-based ordering.
402 SORTABLE_FIELDS: frozenset[str] = frozenset(
403 {
404 "created",
405 "added",
406 "modified",
407 "archive_serial_number",
408 "page_count",
409 "num_notes",
410 },
411 )
413 def __init__(self, path: Path):
414 self._path = path
415 self._raw_index: tantivy.Index | None = None
416 self._raw_schema: tantivy.Schema | None = None
418 @property
419 def _index(self) -> tantivy.Index:
420 assert self._raw_index is not None, "Index not open; call open() first"
421 return self._raw_index
423 @property
424 def _schema(self) -> tantivy.Schema:
425 assert self._raw_schema is not None, "Schema not open; call open() first"
426 return self._raw_schema
428 def open(self) -> None:
429 """
430 Open or rebuild the index as needed.
432 Checks if rebuilding is needed due to schema version or language
433 changes. Registers custom tokenizers after opening.
434 Safe to call multiple times - subsequent calls are no-ops.
435 """
436 if self._raw_index is not None: 436 ↛ 437line 436 didn't jump to line 437 because the condition on line 436 was never true
437 return # pragma: no cover
438 self._raw_index = open_or_rebuild_index(self._path)
439 register_tokenizers(self._raw_index, settings.SEARCH_LANGUAGE)
440 self._raw_schema = self._raw_index.schema
442 def close(self) -> None:
443 """
444 Close the index and release resources.
446 Safe to call multiple times - subsequent calls are no-ops.
447 """
448 self._raw_index = None
449 self._raw_schema = None
451 def _ensure_open(self) -> None:
452 """Ensure the index is open before operations."""
453 if self._raw_index is None: 453 ↛ 454line 453 didn't jump to line 454 because the condition on line 453 was never true
454 self.open() # pragma: no cover
456 def _parse_query(
457 self,
458 query: str,
459 search_mode: SearchMode,
460 ) -> tantivy.Query:
461 """Parse a user query string into a Tantivy Query object."""
462 tz = get_current_timezone()
463 query = normalize_search_text(query)
464 if search_mode is SearchMode.TEXT:
465 return parse_simple_text_query(self._index, query)
466 elif search_mode is SearchMode.TITLE:
467 return parse_simple_title_query(self._index, query)
468 else:
469 return parse_user_query(self._index, query, tz)
471 def _apply_permission_filter(
472 self,
473 query: tantivy.Query,
474 user: AbstractUser | None,
475 ) -> tantivy.Query:
476 """Wrap a query with a permission filter if the user is not a superuser."""
477 if user is not None: 477 ↛ 478line 477 didn't jump to line 478 because the condition on line 477 was never true
478 permission_filter = self._build_permission_filter(user)
479 return tantivy.Query.boolean_query(
480 [
481 (tantivy.Occur.Must, query),
482 (tantivy.Occur.Must, permission_filter),
483 ],
484 )
485 return query
487 def _build_permission_filter(self, user: AbstractUser) -> tantivy.Query:
488 """Build a filter using the user's current group memberships."""
489 group_ids = user.groups.values_list("pk", flat=True)
490 return build_permission_filter(
491 self._schema,
492 user,
493 viewer_group_ids=group_ids,
494 )
496 def _build_tantivy_doc(
497 self,
498 document: Document,
499 viewer_ids: list[int] | None = None,
500 viewer_group_ids: list[int] | None = None,
501 ) -> tantivy.Document:
502 """Build a tantivy Document from a Django Document instance.
504 A root document is indexed with its effective content, i.e. the newest
505 version's OCR text, so it is never indexed with its own outdated text.
506 Annotate the queryset with ``annotate_effective_content`` when indexing
507 more than a couple of documents, to resolve that without a query each.
508 """
509 from guardian.shortcuts import get_groups_with_perms
510 from guardian.shortcuts import get_users_with_perms
512 # Every searchable string is normalized on the way in, and every
513 # query string on the way out (_parse_query), so the two agree on
514 # how a composed character is spelled. See normalize_search_text.
515 content = normalize_search_text(document.get_effective_content() or "")
516 title = normalize_search_text(document.title)
518 doc = tantivy.Document()
520 # Basic fields
521 doc.add_unsigned("id", document.pk)
522 doc.add_text("checksum", document.checksum)
523 doc.add_text("title", title)
524 doc.add_text("title_sort", title)
525 doc.add_text("simple_title", title)
526 doc.add_text("content", content)
527 doc.add_text("simple_content", content)
528 # Bigram (character-ngram) fields exist for CJK substring search,
529 # no need to bloat the bigram index with latin characters.
530 if cjk_title := extract_cjk_text(title):
531 doc.add_text("bigram_title", cjk_title)
532 if content and (cjk_content := extract_cjk_text(content)):
533 doc.add_text("bigram_content", cjk_content)
535 # Original filename - only add if not None/empty
536 if document.original_filename:
537 doc.add_text(
538 "original_filename",
539 normalize_search_text(document.original_filename),
540 )
542 # Correspondent
543 if document.correspondent:
544 correspondent = normalize_search_text(document.correspondent.name)
545 doc.add_text("correspondent", correspondent)
546 doc.add_text("correspondent_sort", correspondent)
547 if cjk_corr := extract_cjk_text(correspondent):
548 doc.add_text("bigram_correspondent", cjk_corr)
550 # Document type
551 if document.document_type:
552 document_type = normalize_search_text(document.document_type.name)
553 doc.add_text("document_type", document_type)
554 doc.add_text("type_sort", document_type)
555 if cjk_type := extract_cjk_text(document_type):
556 doc.add_text("bigram_document_type", cjk_type)
558 # Storage path
559 if document.storage_path:
560 doc.add_text(
561 "storage_path",
562 normalize_search_text(document.storage_path.name),
563 )
565 # Tags — collect names for autocomplete in the same pass
566 tag_names: list[str] = []
567 for tag in document.tags.all():
568 tag_name = normalize_search_text(tag.name)
569 doc.add_text("tag", tag_name)
570 if cjk_tag := extract_cjk_text(tag_name):
571 doc.add_text("bigram_tag", cjk_tag)
572 tag_names.append(tag_name)
574 # Notes — JSON for structured queries (notes.user:alice, notes.note:text).
575 # notes_text is a plain-text companion for snippet/highlight generation;
576 # tantivy's SnippetGenerator does not support JSON fields. It is not in
577 # _DEFAULT_SEARCH_FIELDS, so an unqualified query never searches it: a
578 # note matches through the JSON field or not at all.
579 num_notes = 0
580 note_texts: list[str] = []
581 for note in document.notes.all():
582 num_notes += 1
583 note_text = normalize_search_text(note.note)
584 doc.add_json(
585 "notes",
586 {
587 "note": note_text,
588 "user": (
589 normalize_search_text(note.user.username) if note.user else None
590 ),
591 },
592 )
593 note_texts.append(note_text)
594 if note_texts:
595 doc.add_text("notes_text", " ".join(note_texts))
597 # Custom fields: JSON for structured queries (custom_fields.name:x,
598 # custom_fields.value:y). There is no companion text field here, unlike
599 # notes: custom field values are reachable only through the JSON field.
600 for cfi in document.custom_fields.all():
601 search_value = cfi.value_for_search
602 # Skip fields where there is no value yet
603 if search_value is None:
604 continue
605 doc.add_json(
606 "custom_fields",
607 {
608 "name": normalize_search_text(cfi.field.name),
609 "value": normalize_search_text(search_value),
610 },
611 )
613 # Barcodes: JSON field like custom_fields, only filled when stored
614 for barcode in document.get_effective_barcodes():
615 doc.add_json(
616 "barcodes",
617 {
618 "value": normalize_search_text(barcode.value),
619 "format": normalize_search_text(barcode.format),
620 },
621 )
623 # Dates
624 created_date = datetime(
625 document.created.year,
626 document.created.month,
627 document.created.day,
628 tzinfo=UTC,
629 )
630 doc.add_date("created", created_date)
631 doc.add_date("modified", document.modified)
632 doc.add_date("added", document.added)
634 if document.archive_serial_number is not None:
635 doc.add_unsigned("asn", document.archive_serial_number)
637 if document.page_count is not None:
638 doc.add_unsigned("page_count", document.page_count)
640 doc.add_unsigned("num_notes", num_notes)
642 # Owner
643 if document.owner_id:
644 doc.add_unsigned("owner_id", document.owner_id)
646 # Viewers with permission
647 if viewer_ids is None:
648 users_with_perms = get_users_with_perms(
649 document,
650 only_with_perms_in=["view_document"],
651 with_group_users=False,
652 )
653 viewer_ids = list(
654 cast("QuerySet[User]", users_with_perms).values_list("id", flat=True),
655 )
656 for viewer_id in viewer_ids:
657 doc.add_unsigned("viewer_id", viewer_id)
658 if viewer_group_ids is None:
659 groups_with_perms = get_groups_with_perms(
660 document,
661 only_with_perms_in=["view_document"],
662 )
663 viewer_group_ids = list(
664 cast("QuerySet[Group]", groups_with_perms).values_list(
665 "id",
666 flat=True,
667 ),
668 )
669 for viewer_group_id in viewer_group_ids:
670 doc.add_unsigned("viewer_group_id", viewer_group_id)
672 # Autocomplete words
673 text_sources = [title, content]
674 if document.correspondent:
675 text_sources.append(correspondent)
676 if document.document_type:
677 text_sources.append(document_type)
678 text_sources.extend(tag_names)
680 for word in sorted(_extract_autocomplete_words(text_sources)):
681 doc.add_text("autocomplete_word", word)
683 return doc
685 def add_or_update(self, document: Document) -> None:
686 """
687 Add or update a single document with file locking.
689 Convenience method for single-document updates. For bulk operations,
690 use batch_update() context manager for better performance.
692 On lock exhaustion after all retry attempts, schedules a deferred
693 index_document Celery task and returns normally. Callers will NOT
694 receive a SearchIndexLockError; the index write is deferred silently.
696 Args:
697 document: Django Document instance to index
698 """
699 self._ensure_open()
700 try:
701 with self.batch_update(lock_timeout=_LOCK_TIMEOUT_SECONDS) as batch:
702 batch.add_or_update(document)
703 except SearchIndexLockError:
704 logger.error(
705 "Search index lock exhausted for document %d after %d attempts; "
706 "scheduling deferred index write",
707 document.pk,
708 _LOCK_RETRY_ATTEMPTS,
709 )
710 from documents.tasks import index_document
712 index_document.apply_async(args=[document.pk], countdown=60)
714 def remove(self, doc_id: int) -> None:
715 """
716 Remove a single document from the index with file locking.
718 Convenience method for single-document removal. For bulk operations,
719 use batch_update() context manager for better performance.
721 On lock exhaustion after all retry attempts, schedules a deferred
722 remove_document_from_index Celery task and returns normally.
723 Callers will NOT receive a SearchIndexLockError.
725 Args:
726 doc_id: Primary key of the document to remove
727 """
728 self._ensure_open()
729 try:
730 with self.batch_update(lock_timeout=_LOCK_TIMEOUT_SECONDS) as batch:
731 batch.remove(doc_id)
732 except SearchIndexLockError:
733 logger.error(
734 "Search index lock exhausted for doc_id %d after %d attempts; "
735 "scheduling deferred index removal",
736 doc_id,
737 _LOCK_RETRY_ATTEMPTS,
738 )
739 from documents.tasks import remove_document_from_index
741 remove_document_from_index.apply_async(args=[doc_id], countdown=60)
743 def highlight_hits(
744 self,
745 query: str,
746 doc_ids: list[int],
747 *,
748 search_mode: SearchMode = SearchMode.QUERY,
749 rank_start: int = 1,
750 ) -> list[SearchHit]:
751 """
752 Generate SearchHit dicts with highlights for specific document IDs.
754 Unlike search(), this does not execute a ranked query — it looks up
755 each document by ID and generates snippets against the provided query.
756 Use this when you already know which documents to display (from
757 search_ids + ORM filtering) and just need highlight data.
759 Args:
760 query: The search query (used for snippet generation)
761 doc_ids: Ordered list of document IDs to generate hits for
762 search_mode: Query parsing mode (for building the snippet query)
763 rank_start: Starting rank value (1-based absolute position in the
764 full result set; pass ``page_offset + 1`` for paginated calls)
766 Returns:
767 List of SearchHit dicts in the same order as doc_ids
768 """
769 if not doc_ids: 769 ↛ 772line 769 didn't jump to line 772 because the condition on line 769 was always true
770 return []
772 self._ensure_open()
773 user_query = self._parse_query(query, search_mode)
774 # _parse_query normalizes its own copy; the snippet queries below are
775 # built from the string directly, so normalize it here too.
776 query = normalize_search_text(query)
777 highlight_query = user_query
778 if search_mode is SearchMode.TEXT:
779 try:
780 highlight_query = parse_simple_text_highlight_query(
781 self._index,
782 query,
783 )
784 except ValueError:
785 logger.debug(
786 "Skipping simple text highlight query: token string is not "
787 "valid tantivy query syntax: %r",
788 query,
789 )
791 # For notes_text snippet generation, we need a query that targets the
792 # notes_text field directly. user_query may contain JSON-field terms
793 # (e.g. notes.note:urgent) that the SnippetGenerator cannot resolve
794 # against a text field. Strip field:value prefixes so bare terms like
795 # "urgent" are re-parsed against notes_text, producing highlights even
796 # when the original query used structured syntax.
797 bare_query = re.sub(r"\w[\w.]*:", "", query).strip()
798 try:
799 notes_text_query = (
800 self._index.parse_query(bare_query, ["notes_text"])
801 if bare_query
802 else user_query
803 )
804 except Exception:
805 notes_text_query = user_query
807 searcher = self._index.searcher()
809 # Fetch all requested docs in a single search: user_query MUST match
810 # and exactly the requested IDs MUST match (OR of term_queries).
811 id_filter = tantivy.Query.boolean_query(
812 [
813 (
814 tantivy.Occur.Should,
815 tantivy.Query.term_query(self._schema, "id", did),
816 )
817 for did in doc_ids
818 ],
819 )
820 batch_query = tantivy.Query.boolean_query(
821 [
822 (tantivy.Occur.Must, user_query),
823 (tantivy.Occur.Must, id_filter),
824 ],
825 )
826 batch_results = searcher.search(batch_query, limit=len(doc_ids))
828 result_addrs = [addr for _score, addr in batch_results.hits]
829 result_ids = cast("list[int]", searcher.fast_field_values("id", result_addrs))
830 addr_by_id: dict[int, tuple[float, tantivy.DocAddress]] = {
831 doc_id: (score, addr)
832 for (score, addr), doc_id in zip(batch_results.hits, result_ids)
833 }
835 snippet_generator = None
836 notes_snippet_generator = None
837 hits: list[SearchHit] = []
839 for rank, doc_id in enumerate(doc_ids, start=rank_start):
840 if doc_id not in addr_by_id:
841 continue
843 score, doc_address = addr_by_id[doc_id]
844 actual_doc = searcher.doc(doc_address)
845 doc_dict = actual_doc.to_dict()
847 highlights: dict[str, str] = {}
848 try:
849 if snippet_generator is None:
850 snippet_generator = tantivy.SnippetGenerator.create(
851 searcher,
852 highlight_query,
853 self._schema,
854 "content",
855 )
857 content_html = snippet_generator.snippet_from_doc(actual_doc).to_html()
858 if content_html:
859 highlights["content"] = content_html
861 if search_mode is SearchMode.QUERY and "notes_text" in doc_dict:
862 # Use notes_text (plain text) for snippet generation — tantivy's
863 # SnippetGenerator does not support JSON fields.
864 if notes_snippet_generator is None:
865 notes_snippet_generator = tantivy.SnippetGenerator.create(
866 searcher,
867 notes_text_query,
868 self._schema,
869 "notes_text",
870 )
871 notes_html = notes_snippet_generator.snippet_from_doc(
872 actual_doc,
873 ).to_html()
874 if notes_html:
875 highlights["notes"] = notes_html
877 except Exception: # pragma: no cover
878 logger.debug("Failed to generate highlights for doc %s", doc_id)
880 hits.append(
881 SearchHit(
882 id=doc_id,
883 score=score,
884 rank=rank,
885 highlights=highlights,
886 ),
887 )
889 return hits
891 def search_ids(
892 self,
893 query: str,
894 user: AbstractUser | None,
895 *,
896 sort_field: str | None = None,
897 sort_reverse: bool = False,
898 search_mode: SearchMode = SearchMode.QUERY,
899 limit: int | None = None,
900 ) -> list[int]:
901 """
902 Return document IDs matching a query — no highlights or scores.
904 This is the lightweight companion to search(). Use it when you need the
905 full set of matching IDs (e.g. for ``selection_data``) but don't need
906 scores, ranks, or highlights.
908 Args:
909 query: User's search query
910 user: User for permission filtering (None for superuser/no filtering)
911 sort_field: Field to sort by (None for relevance ranking)
912 sort_reverse: Whether to reverse the sort order
913 search_mode: Query parsing mode (QUERY, TEXT, or TITLE)
914 limit: Maximum number of IDs to return (None = all matching docs)
916 Returns:
917 List of document IDs in the requested order
918 """
919 self._ensure_open()
920 user_query = self._parse_query(query, search_mode)
921 final_query = self._apply_permission_filter(user_query, user)
923 searcher = self._index.searcher()
924 effective_limit = limit if limit is not None else searcher.num_docs
925 if effective_limit <= 0:
926 return []
928 if sort_field and sort_field in self.SORT_FIELD_MAP: 928 ↛ 929line 928 didn't jump to line 929 because the condition on line 928 was never true
929 mapped_field = self.SORT_FIELD_MAP[sort_field]
930 results = searcher.search(
931 final_query,
932 limit=effective_limit,
933 order_by_field=mapped_field,
934 order=tantivy.Order.Desc if sort_reverse else tantivy.Order.Asc,
935 )
936 all_hits = [(hit[1],) for hit in results.hits]
937 else:
938 results = searcher.search(final_query, limit=effective_limit)
939 all_hits = [(hit[1], hit[0]) for hit in results.hits]
941 # Normalize scores and apply threshold (relevance search only)
942 if all_hits: 942 ↛ 943line 942 didn't jump to line 943 because the condition on line 942 was never true
943 max_score = max(hit[1] for hit in all_hits) or 1.0
944 all_hits = [(hit[0], hit[1] / max_score) for hit in all_hits]
946 threshold = settings.ADVANCED_FUZZY_SEARCH_THRESHOLD
947 if threshold is not None: 947 ↛ 948line 947 didn't jump to line 948 because the condition on line 947 was never true
948 all_hits = [hit for hit in all_hits if hit[1] >= threshold]
950 return cast(
951 "list[int]",
952 searcher.fast_field_values("id", [doc_addr for doc_addr, *_ in all_hits]),
953 )
955 def autocomplete(
956 self,
957 term: str,
958 limit: int,
959 user: AbstractUser | None = None,
960 ) -> list[str]:
961 """
962 Get autocomplete suggestions for search queries.
964 Returns words that start with the given term prefix, ranked by document
965 frequency (how many documents contain each word). Optionally filters
966 results to only words from documents visible to the specified user.
968 NOTE: This is the hottest search path (called per keystroke).
969 A future improvement would be to cache results in Redis, keyed by
970 (prefix, user_id), and invalidate on index writes.
972 Args:
973 term: Prefix to match against autocomplete words
974 limit: Maximum number of suggestions to return
975 user: User for permission filtering (None for no filtering)
977 Returns:
978 List of word suggestions ordered by frequency, then alphabetically
979 """
980 self._ensure_open()
981 normalized_term = ascii_fold(term.lower())
982 if not normalized_term:
983 return []
985 searcher = self._index.searcher()
987 permission_query = None
988 # Intersect with permission filter so autocomplete words from
989 # invisible documents don't leak to other users.
990 if user is not None and not user.is_superuser: 990 ↛ 991line 990 didn't jump to line 991 because the condition on line 990 was never true
991 permission_query = self._build_permission_filter(user)
993 matches = searcher.terms_with_prefix(
994 "autocomplete_word",
995 normalized_term,
996 permission_query,
997 limit,
998 )
1000 return [x[0] for x in matches]
1002 def more_like_this_ids(
1003 self,
1004 doc_id: int,
1005 user: AbstractUser | None,
1006 *,
1007 limit: int | None = None,
1008 ) -> list[int]:
1009 """
1010 Return IDs of documents similar to the given document — no highlights.
1012 Lightweight companion to more_like_this(). The original document is
1013 excluded from results.
1015 Args:
1016 doc_id: Primary key of the reference document
1017 user: User for permission filtering (None for no filtering)
1018 limit: Maximum number of IDs to return (None = all matching docs)
1020 Returns:
1021 List of similar document IDs (excluding the original)
1022 """
1023 self._ensure_open()
1024 searcher = self._index.searcher()
1026 id_query = tantivy.Query.term_query(self._schema, "id", doc_id)
1027 results = searcher.search(id_query, limit=1)
1029 if not results.hits:
1030 return []
1032 doc_address = results.hits[0][1]
1033 mlt_query = tantivy.Query.more_like_this_query(
1034 doc_address,
1035 min_doc_frequency=1,
1036 max_doc_frequency=None,
1037 min_term_frequency=1,
1038 max_query_terms=12,
1039 min_word_length=None,
1040 max_word_length=None,
1041 boost_factor=None,
1042 )
1044 final_query = self._apply_permission_filter(mlt_query, user)
1046 effective_limit = limit if limit is not None else searcher.num_docs
1047 try:
1048 # Fetch one extra to account for excluding the original document
1049 results = searcher.search(final_query, limit=effective_limit + 1)
1050 except BaseException: # pragma: no cover
1051 # Tantivy 0.26 panics in BM25 idf scoring when the index holds
1052 # soft-deleted documents (doc_freq can exceed the alive doc count),
1053 # which only surfaces for the More Like This query. The panic crosses
1054 # the pyo3 boundary as a `pyo3_runtime.PanicException` — a
1055 # BaseException, not an Exception — so catch BaseException and degrade
1056 # to "no similar documents" instead of bubbling a 500 to the client.
1057 # Fixed upstream: https://github.com/quickwit-oss/tantivy/pull/2964
1058 # Remove once the bundled tantivy includes that fix.
1059 logger.warning(
1060 "More Like This scoring panicked (likely stale tantivy segment "
1061 "stats after deletions); returning no results. A search index "
1062 "reindex will rebuild consistent statistics.",
1063 )
1064 return []
1066 addrs = [addr for _score, addr in results.hits]
1067 all_ids = cast("list[int]", searcher.fast_field_values("id", addrs))
1068 ids = [rid for rid in all_ids if rid != doc_id]
1069 return ids[:limit] if limit is not None else ids
1071 def batch_update(self, lock_timeout: float = 30.0) -> WriteBatch:
1072 """
1073 Get a batch context manager for bulk index operations.
1075 Use this for efficient bulk document updates/deletions. All operations
1076 within the batch are committed atomically at the end of the context.
1078 Args:
1079 lock_timeout: Seconds to wait for file lock acquisition
1081 Returns:
1082 WriteBatch context manager
1084 Raises:
1085 SearchIndexLockError: If lock cannot be acquired within timeout
1086 """
1087 self._ensure_open()
1088 return WriteBatch(self, lock_timeout)
1090 def rebuild(
1091 self,
1092 documents: QuerySet[Document],
1093 iter_wrapper: IterWrapper[tuple[Document, ViewerGrant]] = identity,
1094 writer_heap_bytes: int = 512_000_000,
1095 ) -> None:
1096 """
1097 Rebuild the entire search index from scratch.
1099 Wipes the existing index and re-indexes all provided documents.
1100 On failure, restores the previous index state to keep the backend usable.
1102 Args:
1103 documents: QuerySet of Document instances to index
1104 iter_wrapper: Optional wrapper function for progress tracking
1105 (e.g., progress bar). Wraps an iterable of
1106 ``(document, (viewer_ids, viewer_group_ids))`` pairs and should yield
1107 each unchanged, advancing one step per document.
1108 writer_heap_bytes: Tantivy writer memory budget (split across the
1109 writer's threads). Larger values buffer more docs in RAM before
1110 flushing a segment, deferring merge work; they do not avoid it.
1111 """
1112 wipe_index(self._path)
1113 new_index = tantivy.Index(build_schema(), path=str(self._path))
1114 _write_sentinels(self._path)
1115 register_tokenizers(new_index, settings.SEARCH_LANGUAGE)
1117 # Point instance at the new index so _build_tantivy_doc uses it
1118 old_index, old_schema = self._raw_index, self._raw_schema
1119 self._raw_index = new_index
1120 self._raw_schema = new_index.schema
1121 # Stream documents one-by-one (so the progress bar advances per
1122 # document) while fetching viewer permissions one SQL query per chunk.
1123 # The stream is Sized, so iter_wrapper can still discover the total.
1124 documents_stream = _DocumentViewerStream(documents, chunk_size=1000)
1125 try:
1126 writer = new_index.writer(heap_size=writer_heap_bytes)
1127 for document, (viewer_ids, viewer_group_ids) in iter_wrapper(
1128 documents_stream,
1129 ):
1130 doc = self._build_tantivy_doc(
1131 document,
1132 viewer_ids=viewer_ids,
1133 viewer_group_ids=viewer_group_ids,
1134 )
1135 writer.add_document(doc)
1136 writer.commit()
1137 # Wait for background merge threads to finish so all segments are
1138 # fully merged and persisted before the index is considered rebuilt.
1139 writer.wait_merging_threads()
1140 new_index.reload()
1141 except BaseException: # pragma: no cover
1142 # Restore old index on failure so the backend remains usable
1143 self._raw_index = old_index
1144 self._raw_schema = old_schema
1145 raise
1148def chunked(iterable, size):
1149 iterator = iter(iterable)
1150 while chunk := list(islice(iterator, size)):
1151 yield chunk
1154_EMPTY_VIEWER_GRANT: Final[ViewerGrant] = ViewerGrant(
1155 viewer_ids=[],
1156 viewer_group_ids=[],
1157)
1160class _DocumentViewerStream(QuerySetStream["Document"]):
1161 """Yield document permission data while batch-loading grants.
1163 Viewer permissions are fetched in batches (see
1164 ``_bulk_get_viewer_permissions``), but documents are yielded individually so a
1165 progress bar wrapped around this stream advances per document rather than
1166 jumping a whole chunk at a time. ``__len__`` (inherited from
1167 ``QuerySetStream``) lets the progress helper still discover the total (it
1168 inspects ``QuerySet``/``Sized``).
1170 The viewer and group ids travel with each document in the yielded pair
1171 rather than through a separate mutable attribute, so the pairing survives
1172 regardless of how ``iter_wrapper`` consumes the stream (buffering,
1173 batching, etc.) — there is no reliance on the caller advancing this
1174 generator in lock-step.
1175 """
1177 def __iter__(self) -> Iterator[tuple[Document, ViewerGrant]]:
1178 # iterator(chunk_size=…) streams from a server-side cursor instead of
1179 # materialising the whole queryset in memory; since Django 4.1 it still
1180 # honours prefetch_related, running the prefetches one batch at a time.
1181 documents = self._queryset.iterator(chunk_size=self._chunk_size)
1182 for chunk in chunked(documents, self._chunk_size):
1183 grants_by_pk = _bulk_get_viewer_permissions([doc.pk for doc in chunk])
1184 for doc in chunk:
1185 yield doc, grants_by_pk.get(doc.pk, _EMPTY_VIEWER_GRANT)
1188def _bulk_get_viewer_permissions(
1189 doc_pks: Sequence[int],
1190) -> dict[int, ViewerGrant]:
1191 """Fetch direct user and group view grants for a batch of documents, keyed by pk.
1193 Group grants remain group IDs in the index so permission checks use the
1194 requesting user's current memberships. Expanding groups to user IDs here
1195 would leave stale access behind after a user is removed from a group.
1196 """
1197 from collections import defaultdict
1199 from django.contrib.contenttypes.models import ContentType
1200 from guardian.models import GroupObjectPermission
1201 from guardian.models import UserObjectPermission
1203 from documents.models import Document
1205 # get_for_model is cached by Django, so this costs at most one query total.
1206 ct = ContentType.objects.get_for_model(Document)
1207 str_pks = [str(pk) for pk in doc_pks]
1209 viewer_map: dict[int, set[int]] = defaultdict(set)
1210 viewer_group_map: dict[int, set[int]] = defaultdict(set)
1212 # Fold the permission lookup into the query via a join on codename instead
1213 # of a separate Permission.objects.get(), which would otherwise run once per
1214 # chunk during a full reindex.
1215 user_qs = UserObjectPermission.objects.filter(
1216 content_type=ct,
1217 permission__content_type=ct,
1218 permission__codename="view_document",
1219 object_pk__in=str_pks,
1220 ).values_list("object_pk", "user_id")
1221 for object_pk, user_id in user_qs:
1222 viewer_map[int(object_pk)].add(user_id)
1224 group_qs = GroupObjectPermission.objects.filter(
1225 content_type=ct,
1226 permission__content_type=ct,
1227 permission__codename="view_document",
1228 object_pk__in=str_pks,
1229 ).values_list("object_pk", "group_id")
1230 for object_pk, group_id in group_qs:
1231 viewer_group_map[int(object_pk)].add(group_id)
1233 return {
1234 object_pk: ViewerGrant(
1235 viewer_ids=list(viewer_map.get(object_pk, ())),
1236 viewer_group_ids=list(viewer_group_map.get(object_pk, ())),
1237 )
1238 for object_pk in viewer_map.keys() | viewer_group_map.keys()
1239 }
1242# Module-level singleton with proper thread safety
1243_backend: TantivyBackend | None = None
1244_backend_path: Path | None = None # tracks which INDEX_DIR the singleton uses
1245_backend_lock = threading.RLock()
1248def get_backend() -> TantivyBackend:
1249 """
1250 Get the global backend instance with thread safety.
1252 Returns a singleton TantivyBackend instance, automatically reinitializing
1253 when settings.INDEX_DIR changes. This ensures proper test isolation when
1254 using pytest-xdist or @override_settings that change the index directory.
1256 Returns:
1257 Thread-safe singleton TantivyBackend instance
1258 """
1259 global _backend, _backend_path
1261 current_path: Path = settings.INDEX_DIR
1263 # Fast path: backend is initialized and path hasn't changed (no lock needed)
1264 if _backend is not None and _backend_path == current_path:
1265 return _backend
1267 # Slow path: first call, or INDEX_DIR changed between calls
1268 with _backend_lock:
1269 # Double-check after acquiring lock — another thread may have beaten us
1270 if _backend is not None and _backend_path == current_path: 1270 ↛ 1271line 1270 didn't jump to line 1271 because the condition on line 1270 was never true
1271 return _backend # pragma: no cover
1273 if _backend is not None: 1273 ↛ 1274line 1273 didn't jump to line 1274 because the condition on line 1273 was never true
1274 _backend.close()
1276 _backend = TantivyBackend(path=current_path)
1277 _backend.open()
1278 _backend_path = current_path
1280 return _backend
1283def reset_backend() -> None:
1284 """
1285 Reset the global backend instance with thread safety.
1287 Forces creation of a new backend instance on the next get_backend() call.
1288 Used for test isolation and when switching between different index directories.
1289 """
1290 global _backend, _backend_path
1292 with _backend_lock:
1293 if _backend is not None:
1294 _backend.close()
1295 _backend = None
1296 _backend_path = None