mirror of
https://github.com/paperless-ngx/paperless-ngx.git
synced 2026-10-03 14:50:31 +00:00
QUERY-mode searches blended a separate bigram clause in at the top of the query, built from the parsed AST's free-text tokens. Because it sat beside the exact clause rather than inside the query, nothing around a CJK term constrained its bigram match: an exclusion that was one OR branch's own condition could never reach it, so "(東京 AND NOT secret) OR bill" still returned the secret document. Widen each CJK leaf where it sits instead, through emit()'s rewrite_leaf hook, so every AND, NOT, REQUIRE, boost and field restriction around the leaf applies to its bigram match too. Negated leaves are widened on purpose, so "NOT X" excludes exactly what "X" matches.
1294 lines
49 KiB
Python
1294 lines
49 KiB
Python
from __future__ import annotations
|
|
|
|
import logging
|
|
import random
|
|
import re
|
|
import threading
|
|
import time
|
|
from datetime import UTC
|
|
from datetime import datetime
|
|
from enum import StrEnum
|
|
from itertools import islice
|
|
from typing import TYPE_CHECKING
|
|
from typing import Final
|
|
from typing import NamedTuple
|
|
from typing import Self
|
|
from typing import TypedDict
|
|
from typing import TypeVar
|
|
from typing import cast
|
|
|
|
import filelock
|
|
import tantivy
|
|
from django.conf import settings
|
|
from django.utils.timezone import get_current_timezone
|
|
|
|
from documents.search._query import extract_cjk_text
|
|
from documents.search._query import normalize_search_text
|
|
from documents.search._query import parse_simple_text_highlight_query
|
|
from documents.search._query import parse_simple_text_query
|
|
from documents.search._query import parse_simple_title_query
|
|
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 wipe_index
|
|
from documents.search._tokenizer import ascii_fold
|
|
from documents.search._tokenizer import autocomplete_tokens
|
|
from documents.search._tokenizer import register_tokenizers
|
|
from documents.utils import IterWrapper
|
|
from documents.utils import QuerySetStream
|
|
from documents.utils import identity
|
|
|
|
if TYPE_CHECKING:
|
|
from collections.abc import Iterable
|
|
from collections.abc import Iterator
|
|
from collections.abc import Sequence
|
|
from pathlib import Path
|
|
|
|
from django.contrib.auth.models import AbstractUser
|
|
from django.contrib.auth.models import Group
|
|
from django.contrib.auth.models import User
|
|
from django.db.models import QuerySet
|
|
|
|
from documents.models import Document
|
|
|
|
logger = logging.getLogger("paperless.search")
|
|
|
|
_LOCK_TIMEOUT_SECONDS: Final[float] = 10.0 # per-attempt acquire timeout
|
|
_LOCK_RETRY_ATTEMPTS: Final[int] = 4 # total attempts (1 initial + 3 retries)
|
|
_LOCK_BACKOFF_BASE: Final[float] = 1.0 # seconds
|
|
_LOCK_BACKOFF_CAP: Final[float] = 10.0 # seconds
|
|
|
|
T = TypeVar("T")
|
|
|
|
|
|
class ViewerGrant(NamedTuple):
|
|
"""Direct user and group view grants for a single document.
|
|
|
|
Named fields (rather than a bare 2-tuple) so ``viewer_ids`` and
|
|
``viewer_group_ids`` can't be silently transposed at a call site — both
|
|
are ``list[int]``, so a positional swap would type-check cleanly.
|
|
"""
|
|
|
|
viewer_ids: list[int]
|
|
viewer_group_ids: list[int]
|
|
|
|
|
|
class SearchMode(StrEnum):
|
|
QUERY = "query"
|
|
TEXT = "text"
|
|
TITLE = "title"
|
|
|
|
|
|
def _extract_autocomplete_words(text_sources: list[str]) -> set[str]:
|
|
"""Extract and normalize words for autocomplete.
|
|
|
|
Tokenizes with Tantivy's simple analyzer (simple -> lowercase -> ascii_fold)
|
|
so the extracted words match how document content is indexed, and runs the
|
|
whole pass in Rust. This replaces a Python regex scan plus per-token folding,
|
|
which dominated full-reindex CPU time. Tokenizing natively also removes the
|
|
ReDoS exposure of running a regex over untrusted content.
|
|
"""
|
|
words = set()
|
|
for text in text_sources:
|
|
if text:
|
|
words.update(autocomplete_tokens(text))
|
|
return words
|
|
|
|
|
|
class SearchHit(TypedDict):
|
|
"""Type definition for search result hits."""
|
|
|
|
id: int
|
|
score: float
|
|
rank: int
|
|
highlights: dict[str, str]
|
|
|
|
|
|
class TantivyRelevanceList:
|
|
"""
|
|
DRF-compatible list wrapper for Tantivy search results.
|
|
|
|
Holds a lightweight ordered list of IDs (for pagination count and
|
|
``selection_data``) together with a small page of rich ``SearchHit``
|
|
dicts (for serialization). DRF's ``PageNumberPagination`` calls
|
|
``__len__`` to compute the total page count and ``__getitem__`` to
|
|
slice the displayed page.
|
|
|
|
Args:
|
|
ordered_ids: All matching document IDs in display order.
|
|
page_hits: Rich SearchHit dicts for the requested DRF page only.
|
|
page_offset: Index into *ordered_ids* where *page_hits* starts.
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
ordered_ids: list[int],
|
|
page_hits: list[SearchHit],
|
|
page_offset: int = 0,
|
|
) -> None:
|
|
self._ordered_ids = ordered_ids
|
|
self._page_hits = page_hits
|
|
self._page_offset = page_offset
|
|
|
|
def __len__(self) -> int:
|
|
return len(self._ordered_ids)
|
|
|
|
def __getitem__(self, key: int | slice) -> SearchHit | list[SearchHit]:
|
|
if isinstance(key, int):
|
|
idx = key if key >= 0 else len(self._ordered_ids) + key
|
|
if self._page_offset <= idx < self._page_offset + len(self._page_hits):
|
|
return self._page_hits[idx - self._page_offset]
|
|
return SearchHit(
|
|
id=self._ordered_ids[key],
|
|
score=0.0,
|
|
rank=idx + 1,
|
|
highlights={},
|
|
)
|
|
start = key.start or 0
|
|
stop = key.stop or len(self._ordered_ids)
|
|
# DRF slices to extract the current page. If the slice aligns
|
|
# with our pre-fetched page_hits, return them directly.
|
|
# We only check start — DRF always slices with stop=start+page_size,
|
|
# which exceeds page_hits length on the last page.
|
|
if start == self._page_offset:
|
|
return self._page_hits[: stop - start]
|
|
# Fallback: return stub dicts (no highlights).
|
|
return [
|
|
SearchHit(id=doc_id, score=0.0, rank=start + i + 1, highlights={})
|
|
for i, doc_id in enumerate(self._ordered_ids[key])
|
|
]
|
|
|
|
def get_all_ids(self) -> list[int]:
|
|
"""Return all matching document IDs in display order."""
|
|
return self._ordered_ids
|
|
|
|
|
|
class SearchIndexLockError(Exception):
|
|
"""Raised when the search index file lock cannot be acquired within the timeout."""
|
|
|
|
|
|
class WriteBatch:
|
|
"""
|
|
Context manager for bulk index operations with file locking.
|
|
|
|
Provides transactional batch updates to the search index with proper
|
|
concurrency control via file locking. All operations within the batch
|
|
are committed atomically or rolled back on exception.
|
|
|
|
Usage:
|
|
with backend.batch_update() as batch:
|
|
batch.add_or_update(document)
|
|
batch.remove(doc_id)
|
|
"""
|
|
|
|
def __init__(self, backend: TantivyBackend, lock_timeout: float):
|
|
self._backend = backend
|
|
self._lock_timeout = lock_timeout
|
|
self._raw_writer: tantivy.IndexWriter | None = None
|
|
self._lock = None
|
|
|
|
@property
|
|
def _writer(self) -> tantivy.IndexWriter:
|
|
assert self._raw_writer is not None, (
|
|
"WriteBatch not entered; use as context manager"
|
|
)
|
|
return self._raw_writer
|
|
|
|
def __enter__(self) -> Self:
|
|
if self._backend._path is not None:
|
|
lock_path = self._backend._path / ".tantivy.lock"
|
|
self._lock = filelock.FileLock(str(lock_path))
|
|
for attempt in range(_LOCK_RETRY_ATTEMPTS):
|
|
try:
|
|
self._lock.acquire(timeout=self._lock_timeout)
|
|
break
|
|
except filelock.Timeout:
|
|
if attempt == _LOCK_RETRY_ATTEMPTS - 1:
|
|
raise SearchIndexLockError(
|
|
f"Could not acquire index lock after {_LOCK_RETRY_ATTEMPTS} "
|
|
f"attempts (timeout={self._lock_timeout}s each)",
|
|
)
|
|
sleep_s = random.uniform(
|
|
0,
|
|
min(_LOCK_BACKOFF_CAP, _LOCK_BACKOFF_BASE * (2**attempt)),
|
|
)
|
|
logger.debug(
|
|
"Index lock contention; retrying in %.2fs (attempt %d/%d)",
|
|
sleep_s,
|
|
attempt + 1,
|
|
_LOCK_RETRY_ATTEMPTS,
|
|
)
|
|
time.sleep(sleep_s)
|
|
|
|
# Open a fresh Index (and thus a fresh Tantivy ManagedDirectory)
|
|
# for the write, rather than reusing the process-local cached
|
|
# index. ManagedDirectory loads its GC bookkeeping (.managed.json)
|
|
# once, at construction, and never re-reads it; paperless runs
|
|
# several long-lived processes (Granian workers, Celery workers)
|
|
# that take turns writing under the file lock above. A cached,
|
|
# long-lived writer index would carry a stale managed-files view
|
|
# and, on commit, overwrite .managed.json with that stale view -
|
|
# permanently losing track of segment files other processes
|
|
# registered in the meantime, so they can never be garbage
|
|
# collected. Reopening fresh here always picks up the current
|
|
# on-disk state. The long-lived self._backend._index is used for
|
|
# reads only and is reloaded (not reopened) after commit below.
|
|
write_index = tantivy.Index(
|
|
build_schema(),
|
|
path=str(self._backend._path),
|
|
)
|
|
register_tokenizers(write_index, settings.SEARCH_LANGUAGE)
|
|
self._raw_writer = write_index.writer()
|
|
else:
|
|
self._raw_writer = self._backend._index.writer()
|
|
return self
|
|
|
|
def __exit__(self, exc_type, exc_val, exc_tb):
|
|
try:
|
|
if exc_type is None:
|
|
self._writer.commit()
|
|
# Wait for background merge threads to finish before releasing
|
|
# the file lock so the next writer doesn't race against an
|
|
# in-progress merge on the same index files.
|
|
self._writer.wait_merging_threads()
|
|
self._backend._index.reload()
|
|
finally:
|
|
# Always release the writer (and Tantivy's internal writer lock),
|
|
# even if commit/merge/reload raised, so the next batch can acquire
|
|
# a writer instead of failing with LockBusy. An uncommitted writer
|
|
# is simply discarded.
|
|
if self._raw_writer is not None:
|
|
del self._raw_writer
|
|
self._raw_writer = None
|
|
if self._lock is not None:
|
|
self._lock.release()
|
|
|
|
def add_or_update(self, document: Document) -> None:
|
|
"""
|
|
Add or update a document in the batch.
|
|
|
|
Implements upsert behavior by deleting any existing document with the same ID
|
|
and adding the new version. This ensures stale document data (e.g., after
|
|
permission changes) doesn't persist in the index.
|
|
|
|
Args:
|
|
document: Django Document instance to index
|
|
"""
|
|
self.remove(document.pk)
|
|
doc = self._backend._build_tantivy_doc(document)
|
|
self._writer.add_document(doc)
|
|
|
|
def remove(self, doc_id: int) -> None:
|
|
"""Remove a document from the batch by its primary key."""
|
|
self._writer.delete_documents_by_query(
|
|
tantivy.Query.term_query(self._backend._schema, "id", doc_id),
|
|
)
|
|
|
|
def add_or_update_ids(self, ids: Sequence[int]) -> None:
|
|
"""
|
|
Add or update multiple documents in the batch by primary key.
|
|
|
|
Unlike calling ``add_or_update()`` once per document, this resolves
|
|
viewer permissions and effective (versioned) content in bulk against
|
|
the ids as a whole, instead of once per document -- see
|
|
``_DocumentViewerStream`` and ``annotate_effective_content``. Use
|
|
this whenever more than one document is being written in the same
|
|
batch.
|
|
|
|
An id with no matching document (e.g. deleted between the caller
|
|
collecting ids and the batch running) is silently skipped, matching
|
|
``add_or_update()``'s existing single-document deferred-task behavior
|
|
rather than erroring or leaving a stale index entry.
|
|
|
|
Args:
|
|
ids: Primary keys of Document instances to index
|
|
"""
|
|
from documents.models import Document
|
|
from documents.versioning import annotate_effective_content
|
|
|
|
ids = list(ids)
|
|
if not ids:
|
|
return
|
|
|
|
queryset = annotate_effective_content(
|
|
Document.objects.filter(pk__in=ids)
|
|
.select_related("correspondent", "document_type", "storage_path", "owner")
|
|
.prefetch_related("tags", "notes__user", "custom_fields__field"),
|
|
)
|
|
for document, grant in _DocumentViewerStream(queryset, chunk_size=1000):
|
|
self.remove(document.pk)
|
|
doc = self._backend._build_tantivy_doc(
|
|
document,
|
|
viewer_ids=grant.viewer_ids,
|
|
viewer_group_ids=grant.viewer_group_ids,
|
|
)
|
|
self._writer.add_document(doc)
|
|
|
|
|
|
def build_permission_filter(
|
|
schema: tantivy.Schema,
|
|
user: AbstractUser,
|
|
viewer_group_ids: Iterable[int] = (),
|
|
) -> tantivy.Query:
|
|
"""
|
|
Build a query filter for user document permissions.
|
|
|
|
Creates a query that matches only documents visible to the specified user
|
|
according to paperless-ngx permission rules:
|
|
- Public documents (no owner) are visible to all users
|
|
- Private documents are visible to their owner
|
|
- Documents explicitly shared with the user are visible
|
|
- Documents shared with one of the user's current groups are visible
|
|
|
|
Args:
|
|
schema: Tantivy schema for field validation
|
|
user: User to check permissions for
|
|
viewer_group_ids: Current group memberships for the user
|
|
|
|
Returns:
|
|
Tantivy query that filters results to visible documents
|
|
"""
|
|
owner_any = tantivy.Query.exists_query("owner_id")
|
|
no_owner = tantivy.Query.boolean_query(
|
|
[
|
|
(tantivy.Occur.Must, tantivy.Query.all_query()),
|
|
(tantivy.Occur.MustNot, owner_any),
|
|
],
|
|
)
|
|
owned = tantivy.Query.term_query(schema, "owner_id", user.pk)
|
|
shared = tantivy.Query.term_query(schema, "viewer_id", user.pk)
|
|
group_shared = [
|
|
tantivy.Query.term_query(schema, "viewer_group_id", group_id)
|
|
for group_id in viewer_group_ids
|
|
]
|
|
return tantivy.Query.disjunction_max_query(
|
|
[no_owner, owned, shared, *group_shared],
|
|
)
|
|
|
|
|
|
class TantivyBackend:
|
|
"""
|
|
Tantivy search backend with explicit lifecycle management.
|
|
|
|
Provides full-text search capabilities using the Tantivy search engine.
|
|
Supports in-memory indexes (for testing) and persistent on-disk indexes
|
|
(for production use). Handles document indexing, search queries, autocompletion,
|
|
and "more like this" functionality.
|
|
|
|
The backend manages its own connection lifecycle and can be reset when
|
|
the underlying index directory changes (e.g., during test isolation).
|
|
"""
|
|
|
|
# Maps DRF ordering field names to Tantivy index field names.
|
|
SORT_FIELD_MAP: dict[str, str] = {
|
|
"title": "title_sort",
|
|
"correspondent__name": "correspondent_sort",
|
|
"document_type__name": "type_sort",
|
|
"created": "created",
|
|
"added": "added",
|
|
"modified": "modified",
|
|
"archive_serial_number": "asn",
|
|
"page_count": "page_count",
|
|
"num_notes": "num_notes",
|
|
}
|
|
|
|
# Fields where Tantivy's sort order matches the ORM's sort order.
|
|
# Text-based fields (title, correspondent__name, document_type__name)
|
|
# are excluded because Tantivy's tokenized fast fields produce different
|
|
# ordering than the ORM's collation-based ordering.
|
|
SORTABLE_FIELDS: frozenset[str] = frozenset(
|
|
{
|
|
"created",
|
|
"added",
|
|
"modified",
|
|
"archive_serial_number",
|
|
"page_count",
|
|
"num_notes",
|
|
},
|
|
)
|
|
|
|
def __init__(self, path: Path | None = None):
|
|
# path=None → in-memory index (for tests)
|
|
# path=some_dir → on-disk index (for production)
|
|
self._path = path
|
|
self._raw_index: tantivy.Index | None = None
|
|
self._raw_schema: tantivy.Schema | None = None
|
|
|
|
@property
|
|
def _index(self) -> tantivy.Index:
|
|
assert self._raw_index is not None, "Index not open; call open() first"
|
|
return self._raw_index
|
|
|
|
@property
|
|
def _schema(self) -> tantivy.Schema:
|
|
assert self._raw_schema is not None, "Schema not open; call open() first"
|
|
return self._raw_schema
|
|
|
|
def open(self) -> None:
|
|
"""
|
|
Open or rebuild the index as needed.
|
|
|
|
For disk-based indexes, checks if rebuilding is needed due to schema
|
|
version or language changes. Registers custom tokenizers after opening.
|
|
Safe to call multiple times - subsequent calls are no-ops.
|
|
"""
|
|
if self._raw_index is not None:
|
|
return # pragma: no cover
|
|
if self._path is not None:
|
|
self._raw_index = open_or_rebuild_index(self._path)
|
|
else:
|
|
self._raw_index = tantivy.Index(build_schema())
|
|
register_tokenizers(self._raw_index, settings.SEARCH_LANGUAGE)
|
|
self._raw_schema = self._raw_index.schema
|
|
|
|
def close(self) -> None:
|
|
"""
|
|
Close the index and release resources.
|
|
|
|
Safe to call multiple times - subsequent calls are no-ops.
|
|
"""
|
|
self._raw_index = None
|
|
self._raw_schema = None
|
|
|
|
def _ensure_open(self) -> None:
|
|
"""Ensure the index is open before operations."""
|
|
if self._raw_index is None:
|
|
self.open() # pragma: no cover
|
|
|
|
def _parse_query(
|
|
self,
|
|
query: str,
|
|
search_mode: SearchMode,
|
|
) -> tantivy.Query:
|
|
"""Parse a user query string into a Tantivy Query object."""
|
|
tz = get_current_timezone()
|
|
query = normalize_search_text(query)
|
|
if search_mode is SearchMode.TEXT:
|
|
return parse_simple_text_query(self._index, query)
|
|
elif search_mode is SearchMode.TITLE:
|
|
return parse_simple_title_query(self._index, query)
|
|
else:
|
|
return parse_user_query(self._index, query, tz)
|
|
|
|
def _apply_permission_filter(
|
|
self,
|
|
query: tantivy.Query,
|
|
user: AbstractUser | None,
|
|
) -> tantivy.Query:
|
|
"""Wrap a query with a permission filter if the user is not a superuser."""
|
|
if user is not None:
|
|
permission_filter = self._build_permission_filter(user)
|
|
return tantivy.Query.boolean_query(
|
|
[
|
|
(tantivy.Occur.Must, query),
|
|
(tantivy.Occur.Must, permission_filter),
|
|
],
|
|
)
|
|
return query
|
|
|
|
def _build_permission_filter(self, user: AbstractUser) -> tantivy.Query:
|
|
"""Build a filter using the user's current group memberships."""
|
|
group_ids = user.groups.values_list("pk", flat=True)
|
|
return build_permission_filter(
|
|
self._schema,
|
|
user,
|
|
viewer_group_ids=group_ids,
|
|
)
|
|
|
|
def _build_tantivy_doc(
|
|
self,
|
|
document: Document,
|
|
viewer_ids: list[int] | None = None,
|
|
viewer_group_ids: list[int] | None = None,
|
|
) -> tantivy.Document:
|
|
"""Build a tantivy Document from a Django Document instance.
|
|
|
|
A root document is indexed with its effective content, i.e. the newest
|
|
version's OCR text, so it is never indexed with its own outdated text.
|
|
Annotate the queryset with ``annotate_effective_content`` when indexing
|
|
more than a couple of documents, to resolve that without a query each.
|
|
"""
|
|
from guardian.shortcuts import get_groups_with_perms
|
|
from guardian.shortcuts import get_users_with_perms
|
|
|
|
# Every searchable string is normalized on the way in, and every
|
|
# query string on the way out (_parse_query), so the two agree on
|
|
# how a composed character is spelled. See normalize_search_text.
|
|
content = normalize_search_text(document.get_effective_content() or "")
|
|
title = normalize_search_text(document.title)
|
|
|
|
doc = tantivy.Document()
|
|
|
|
# Basic fields
|
|
doc.add_unsigned("id", document.pk)
|
|
doc.add_text("checksum", document.checksum)
|
|
doc.add_text("title", title)
|
|
doc.add_text("title_sort", title)
|
|
doc.add_text("simple_title", title)
|
|
doc.add_text("content", content)
|
|
doc.add_text("simple_content", content)
|
|
# Bigram (character-ngram) fields exist for CJK substring search,
|
|
# no need to bloat the bigram index with latin characters.
|
|
if cjk_title := extract_cjk_text(title):
|
|
doc.add_text("bigram_title", cjk_title)
|
|
if content and (cjk_content := extract_cjk_text(content)):
|
|
doc.add_text("bigram_content", cjk_content)
|
|
|
|
# Original filename - only add if not None/empty
|
|
if document.original_filename:
|
|
doc.add_text(
|
|
"original_filename",
|
|
normalize_search_text(document.original_filename),
|
|
)
|
|
|
|
# Correspondent
|
|
if document.correspondent:
|
|
correspondent = normalize_search_text(document.correspondent.name)
|
|
doc.add_text("correspondent", correspondent)
|
|
doc.add_text("correspondent_sort", correspondent)
|
|
if cjk_corr := extract_cjk_text(correspondent):
|
|
doc.add_text("bigram_correspondent", cjk_corr)
|
|
|
|
# Document type
|
|
if document.document_type:
|
|
document_type = normalize_search_text(document.document_type.name)
|
|
doc.add_text("document_type", document_type)
|
|
doc.add_text("type_sort", document_type)
|
|
if cjk_type := extract_cjk_text(document_type):
|
|
doc.add_text("bigram_document_type", cjk_type)
|
|
|
|
# Storage path
|
|
if document.storage_path:
|
|
doc.add_text(
|
|
"storage_path",
|
|
normalize_search_text(document.storage_path.name),
|
|
)
|
|
|
|
# Tags — collect names for autocomplete in the same pass
|
|
tag_names: list[str] = []
|
|
for tag in document.tags.all():
|
|
tag_name = normalize_search_text(tag.name)
|
|
doc.add_text("tag", tag_name)
|
|
if cjk_tag := extract_cjk_text(tag_name):
|
|
doc.add_text("bigram_tag", cjk_tag)
|
|
tag_names.append(tag_name)
|
|
|
|
# Notes — JSON for structured queries (notes.user:alice, notes.note:text).
|
|
# notes_text is a plain-text companion for snippet/highlight generation;
|
|
# tantivy's SnippetGenerator does not support JSON fields. It is not in
|
|
# _DEFAULT_SEARCH_FIELDS, so an unqualified query never searches it: a
|
|
# note matches through the JSON field or not at all.
|
|
num_notes = 0
|
|
note_texts: list[str] = []
|
|
for note in document.notes.all():
|
|
num_notes += 1
|
|
note_text = normalize_search_text(note.note)
|
|
doc.add_json(
|
|
"notes",
|
|
{
|
|
"note": note_text,
|
|
"user": (
|
|
normalize_search_text(note.user.username) if note.user else None
|
|
),
|
|
},
|
|
)
|
|
note_texts.append(note_text)
|
|
if note_texts:
|
|
doc.add_text("notes_text", " ".join(note_texts))
|
|
|
|
# Custom fields: JSON for structured queries (custom_fields.name:x,
|
|
# custom_fields.value:y). There is no companion text field here, unlike
|
|
# notes: custom field values are reachable only through the JSON field.
|
|
for cfi in document.custom_fields.all():
|
|
search_value = cfi.value_for_search
|
|
# Skip fields where there is no value yet
|
|
if search_value is None:
|
|
continue
|
|
doc.add_json(
|
|
"custom_fields",
|
|
{
|
|
"name": normalize_search_text(cfi.field.name),
|
|
"value": normalize_search_text(search_value),
|
|
},
|
|
)
|
|
|
|
# Dates
|
|
created_date = datetime(
|
|
document.created.year,
|
|
document.created.month,
|
|
document.created.day,
|
|
tzinfo=UTC,
|
|
)
|
|
doc.add_date("created", created_date)
|
|
doc.add_date("modified", document.modified)
|
|
doc.add_date("added", document.added)
|
|
|
|
if document.archive_serial_number is not None:
|
|
doc.add_unsigned("asn", document.archive_serial_number)
|
|
|
|
if document.page_count is not None:
|
|
doc.add_unsigned("page_count", document.page_count)
|
|
|
|
doc.add_unsigned("num_notes", num_notes)
|
|
|
|
# Owner
|
|
if document.owner_id:
|
|
doc.add_unsigned("owner_id", document.owner_id)
|
|
|
|
# Viewers with permission
|
|
if viewer_ids is None:
|
|
users_with_perms = get_users_with_perms(
|
|
document,
|
|
only_with_perms_in=["view_document"],
|
|
with_group_users=False,
|
|
)
|
|
viewer_ids = list(
|
|
cast("QuerySet[User]", users_with_perms).values_list("id", flat=True),
|
|
)
|
|
for viewer_id in viewer_ids:
|
|
doc.add_unsigned("viewer_id", viewer_id)
|
|
if viewer_group_ids is None:
|
|
groups_with_perms = get_groups_with_perms(
|
|
document,
|
|
only_with_perms_in=["view_document"],
|
|
)
|
|
viewer_group_ids = list(
|
|
cast("QuerySet[Group]", groups_with_perms).values_list(
|
|
"id",
|
|
flat=True,
|
|
),
|
|
)
|
|
for viewer_group_id in viewer_group_ids:
|
|
doc.add_unsigned("viewer_group_id", viewer_group_id)
|
|
|
|
# Autocomplete words
|
|
text_sources = [title, content]
|
|
if document.correspondent:
|
|
text_sources.append(correspondent)
|
|
if document.document_type:
|
|
text_sources.append(document_type)
|
|
text_sources.extend(tag_names)
|
|
|
|
for word in sorted(_extract_autocomplete_words(text_sources)):
|
|
doc.add_text("autocomplete_word", word)
|
|
|
|
return doc
|
|
|
|
def add_or_update(self, document: Document) -> None:
|
|
"""
|
|
Add or update a single document with file locking.
|
|
|
|
Convenience method for single-document updates. For bulk operations,
|
|
use batch_update() context manager for better performance.
|
|
|
|
On lock exhaustion after all retry attempts, schedules a deferred
|
|
index_document Celery task and returns normally. Callers will NOT
|
|
receive a SearchIndexLockError; the index write is deferred silently.
|
|
|
|
Args:
|
|
document: Django Document instance to index
|
|
"""
|
|
self._ensure_open()
|
|
try:
|
|
with self.batch_update(lock_timeout=_LOCK_TIMEOUT_SECONDS) as batch:
|
|
batch.add_or_update(document)
|
|
except SearchIndexLockError:
|
|
logger.error(
|
|
"Search index lock exhausted for document %d after %d attempts; "
|
|
"scheduling deferred index write",
|
|
document.pk,
|
|
_LOCK_RETRY_ATTEMPTS,
|
|
)
|
|
from documents.tasks import index_document
|
|
|
|
index_document.apply_async(args=[document.pk], countdown=60)
|
|
|
|
def remove(self, doc_id: int) -> None:
|
|
"""
|
|
Remove a single document from the index with file locking.
|
|
|
|
Convenience method for single-document removal. For bulk operations,
|
|
use batch_update() context manager for better performance.
|
|
|
|
On lock exhaustion after all retry attempts, schedules a deferred
|
|
remove_document_from_index Celery task and returns normally.
|
|
Callers will NOT receive a SearchIndexLockError.
|
|
|
|
Args:
|
|
doc_id: Primary key of the document to remove
|
|
"""
|
|
self._ensure_open()
|
|
try:
|
|
with self.batch_update(lock_timeout=_LOCK_TIMEOUT_SECONDS) as batch:
|
|
batch.remove(doc_id)
|
|
except SearchIndexLockError:
|
|
logger.error(
|
|
"Search index lock exhausted for doc_id %d after %d attempts; "
|
|
"scheduling deferred index removal",
|
|
doc_id,
|
|
_LOCK_RETRY_ATTEMPTS,
|
|
)
|
|
from documents.tasks import remove_document_from_index
|
|
|
|
remove_document_from_index.apply_async(args=[doc_id], countdown=60)
|
|
|
|
def highlight_hits(
|
|
self,
|
|
query: str,
|
|
doc_ids: list[int],
|
|
*,
|
|
search_mode: SearchMode = SearchMode.QUERY,
|
|
rank_start: int = 1,
|
|
) -> list[SearchHit]:
|
|
"""
|
|
Generate SearchHit dicts with highlights for specific document IDs.
|
|
|
|
Unlike search(), this does not execute a ranked query — it looks up
|
|
each document by ID and generates snippets against the provided query.
|
|
Use this when you already know which documents to display (from
|
|
search_ids + ORM filtering) and just need highlight data.
|
|
|
|
Args:
|
|
query: The search query (used for snippet generation)
|
|
doc_ids: Ordered list of document IDs to generate hits for
|
|
search_mode: Query parsing mode (for building the snippet query)
|
|
rank_start: Starting rank value (1-based absolute position in the
|
|
full result set; pass ``page_offset + 1`` for paginated calls)
|
|
|
|
Returns:
|
|
List of SearchHit dicts in the same order as doc_ids
|
|
"""
|
|
if not doc_ids:
|
|
return []
|
|
|
|
self._ensure_open()
|
|
user_query = self._parse_query(query, search_mode)
|
|
# _parse_query normalizes its own copy; the snippet queries below are
|
|
# built from the string directly, so normalize it here too.
|
|
query = normalize_search_text(query)
|
|
highlight_query = user_query
|
|
if search_mode is SearchMode.TEXT:
|
|
try:
|
|
highlight_query = parse_simple_text_highlight_query(
|
|
self._index,
|
|
query,
|
|
)
|
|
except ValueError:
|
|
logger.debug(
|
|
"Skipping simple text highlight query: token string is not "
|
|
"valid tantivy query syntax: %r",
|
|
query,
|
|
)
|
|
|
|
# For notes_text snippet generation, we need a query that targets the
|
|
# notes_text field directly. user_query may contain JSON-field terms
|
|
# (e.g. notes.note:urgent) that the SnippetGenerator cannot resolve
|
|
# against a text field. Strip field:value prefixes so bare terms like
|
|
# "urgent" are re-parsed against notes_text, producing highlights even
|
|
# when the original query used structured syntax.
|
|
bare_query = re.sub(r"\w[\w.]*:", "", query).strip()
|
|
try:
|
|
notes_text_query = (
|
|
self._index.parse_query(bare_query, ["notes_text"])
|
|
if bare_query
|
|
else user_query
|
|
)
|
|
except Exception:
|
|
notes_text_query = user_query
|
|
|
|
searcher = self._index.searcher()
|
|
|
|
# Fetch all requested docs in a single search: user_query MUST match
|
|
# and exactly the requested IDs MUST match (OR of term_queries).
|
|
id_filter = tantivy.Query.boolean_query(
|
|
[
|
|
(
|
|
tantivy.Occur.Should,
|
|
tantivy.Query.term_query(self._schema, "id", did),
|
|
)
|
|
for did in doc_ids
|
|
],
|
|
)
|
|
batch_query = tantivy.Query.boolean_query(
|
|
[
|
|
(tantivy.Occur.Must, user_query),
|
|
(tantivy.Occur.Must, id_filter),
|
|
],
|
|
)
|
|
batch_results = searcher.search(batch_query, limit=len(doc_ids))
|
|
|
|
result_addrs = [addr for _score, addr in batch_results.hits]
|
|
result_ids = cast("list[int]", searcher.fast_field_values("id", result_addrs))
|
|
addr_by_id: dict[int, tuple[float, tantivy.DocAddress]] = {
|
|
doc_id: (score, addr)
|
|
for (score, addr), doc_id in zip(batch_results.hits, result_ids)
|
|
}
|
|
|
|
snippet_generator = None
|
|
notes_snippet_generator = None
|
|
hits: list[SearchHit] = []
|
|
|
|
for rank, doc_id in enumerate(doc_ids, start=rank_start):
|
|
if doc_id not in addr_by_id:
|
|
continue
|
|
|
|
score, doc_address = addr_by_id[doc_id]
|
|
actual_doc = searcher.doc(doc_address)
|
|
doc_dict = actual_doc.to_dict()
|
|
|
|
highlights: dict[str, str] = {}
|
|
try:
|
|
if snippet_generator is None:
|
|
snippet_generator = tantivy.SnippetGenerator.create(
|
|
searcher,
|
|
highlight_query,
|
|
self._schema,
|
|
"content",
|
|
)
|
|
|
|
content_html = snippet_generator.snippet_from_doc(actual_doc).to_html()
|
|
if content_html:
|
|
highlights["content"] = content_html
|
|
|
|
if search_mode is SearchMode.QUERY and "notes_text" in doc_dict:
|
|
# Use notes_text (plain text) for snippet generation — tantivy's
|
|
# SnippetGenerator does not support JSON fields.
|
|
if notes_snippet_generator is None:
|
|
notes_snippet_generator = tantivy.SnippetGenerator.create(
|
|
searcher,
|
|
notes_text_query,
|
|
self._schema,
|
|
"notes_text",
|
|
)
|
|
notes_html = notes_snippet_generator.snippet_from_doc(
|
|
actual_doc,
|
|
).to_html()
|
|
if notes_html:
|
|
highlights["notes"] = notes_html
|
|
|
|
except Exception: # pragma: no cover
|
|
logger.debug("Failed to generate highlights for doc %s", doc_id)
|
|
|
|
hits.append(
|
|
SearchHit(
|
|
id=doc_id,
|
|
score=score,
|
|
rank=rank,
|
|
highlights=highlights,
|
|
),
|
|
)
|
|
|
|
return hits
|
|
|
|
def search_ids(
|
|
self,
|
|
query: str,
|
|
user: AbstractUser | None,
|
|
*,
|
|
sort_field: str | None = None,
|
|
sort_reverse: bool = False,
|
|
search_mode: SearchMode = SearchMode.QUERY,
|
|
limit: int | None = None,
|
|
) -> list[int]:
|
|
"""
|
|
Return document IDs matching a query — no highlights or scores.
|
|
|
|
This is the lightweight companion to search(). Use it when you need the
|
|
full set of matching IDs (e.g. for ``selection_data``) but don't need
|
|
scores, ranks, or highlights.
|
|
|
|
Args:
|
|
query: User's search query
|
|
user: User for permission filtering (None for superuser/no filtering)
|
|
sort_field: Field to sort by (None for relevance ranking)
|
|
sort_reverse: Whether to reverse the sort order
|
|
search_mode: Query parsing mode (QUERY, TEXT, or TITLE)
|
|
limit: Maximum number of IDs to return (None = all matching docs)
|
|
|
|
Returns:
|
|
List of document IDs in the requested order
|
|
"""
|
|
self._ensure_open()
|
|
user_query = self._parse_query(query, search_mode)
|
|
final_query = self._apply_permission_filter(user_query, user)
|
|
|
|
searcher = self._index.searcher()
|
|
effective_limit = limit if limit is not None else searcher.num_docs
|
|
if effective_limit <= 0:
|
|
return []
|
|
|
|
if sort_field and sort_field in self.SORT_FIELD_MAP:
|
|
mapped_field = self.SORT_FIELD_MAP[sort_field]
|
|
results = searcher.search(
|
|
final_query,
|
|
limit=effective_limit,
|
|
order_by_field=mapped_field,
|
|
order=tantivy.Order.Desc if sort_reverse else tantivy.Order.Asc,
|
|
)
|
|
all_hits = [(hit[1],) for hit in results.hits]
|
|
else:
|
|
results = searcher.search(final_query, limit=effective_limit)
|
|
all_hits = [(hit[1], hit[0]) for hit in results.hits]
|
|
|
|
# Normalize scores and apply threshold (relevance search only)
|
|
if all_hits:
|
|
max_score = max(hit[1] for hit in all_hits) or 1.0
|
|
all_hits = [(hit[0], hit[1] / max_score) for hit in all_hits]
|
|
|
|
threshold = settings.ADVANCED_FUZZY_SEARCH_THRESHOLD
|
|
if threshold is not None:
|
|
all_hits = [hit for hit in all_hits if hit[1] >= threshold]
|
|
|
|
return cast(
|
|
"list[int]",
|
|
searcher.fast_field_values("id", [doc_addr for doc_addr, *_ in all_hits]),
|
|
)
|
|
|
|
def autocomplete(
|
|
self,
|
|
term: str,
|
|
limit: int,
|
|
user: AbstractUser | None = None,
|
|
) -> list[str]:
|
|
"""
|
|
Get autocomplete suggestions for search queries.
|
|
|
|
Returns words that start with the given term prefix, ranked by document
|
|
frequency (how many documents contain each word). Optionally filters
|
|
results to only words from documents visible to the specified user.
|
|
|
|
NOTE: This is the hottest search path (called per keystroke).
|
|
A future improvement would be to cache results in Redis, keyed by
|
|
(prefix, user_id), and invalidate on index writes.
|
|
|
|
Args:
|
|
term: Prefix to match against autocomplete words
|
|
limit: Maximum number of suggestions to return
|
|
user: User for permission filtering (None for no filtering)
|
|
|
|
Returns:
|
|
List of word suggestions ordered by frequency, then alphabetically
|
|
"""
|
|
self._ensure_open()
|
|
normalized_term = ascii_fold(term.lower())
|
|
if not normalized_term:
|
|
return []
|
|
|
|
searcher = self._index.searcher()
|
|
|
|
permission_query = None
|
|
# Intersect with permission filter so autocomplete words from
|
|
# invisible documents don't leak to other users.
|
|
if user is not None and not user.is_superuser:
|
|
permission_query = self._build_permission_filter(user)
|
|
|
|
matches = searcher.terms_with_prefix(
|
|
"autocomplete_word",
|
|
normalized_term,
|
|
permission_query,
|
|
limit,
|
|
)
|
|
|
|
return [x[0] for x in matches]
|
|
|
|
def more_like_this_ids(
|
|
self,
|
|
doc_id: int,
|
|
user: AbstractUser | None,
|
|
*,
|
|
limit: int | None = None,
|
|
) -> list[int]:
|
|
"""
|
|
Return IDs of documents similar to the given document — no highlights.
|
|
|
|
Lightweight companion to more_like_this(). The original document is
|
|
excluded from results.
|
|
|
|
Args:
|
|
doc_id: Primary key of the reference document
|
|
user: User for permission filtering (None for no filtering)
|
|
limit: Maximum number of IDs to return (None = all matching docs)
|
|
|
|
Returns:
|
|
List of similar document IDs (excluding the original)
|
|
"""
|
|
self._ensure_open()
|
|
searcher = self._index.searcher()
|
|
|
|
id_query = tantivy.Query.term_query(self._schema, "id", doc_id)
|
|
results = searcher.search(id_query, limit=1)
|
|
|
|
if not results.hits:
|
|
return []
|
|
|
|
doc_address = results.hits[0][1]
|
|
mlt_query = tantivy.Query.more_like_this_query(
|
|
doc_address,
|
|
min_doc_frequency=1,
|
|
max_doc_frequency=None,
|
|
min_term_frequency=1,
|
|
max_query_terms=12,
|
|
min_word_length=None,
|
|
max_word_length=None,
|
|
boost_factor=None,
|
|
)
|
|
|
|
final_query = self._apply_permission_filter(mlt_query, user)
|
|
|
|
effective_limit = limit if limit is not None else searcher.num_docs
|
|
try:
|
|
# Fetch one extra to account for excluding the original document
|
|
results = searcher.search(final_query, limit=effective_limit + 1)
|
|
except BaseException: # pragma: no cover
|
|
# Tantivy 0.26 panics in BM25 idf scoring when the index holds
|
|
# soft-deleted documents (doc_freq can exceed the alive doc count),
|
|
# which only surfaces for the More Like This query. The panic crosses
|
|
# the pyo3 boundary as a `pyo3_runtime.PanicException` — a
|
|
# BaseException, not an Exception — so catch BaseException and degrade
|
|
# to "no similar documents" instead of bubbling a 500 to the client.
|
|
# Fixed upstream: https://github.com/quickwit-oss/tantivy/pull/2964
|
|
# Remove once the bundled tantivy includes that fix.
|
|
logger.warning(
|
|
"More Like This scoring panicked (likely stale tantivy segment "
|
|
"stats after deletions); returning no results. A search index "
|
|
"reindex will rebuild consistent statistics.",
|
|
)
|
|
return []
|
|
|
|
addrs = [addr for _score, addr in results.hits]
|
|
all_ids = cast("list[int]", searcher.fast_field_values("id", addrs))
|
|
ids = [rid for rid in all_ids if rid != doc_id]
|
|
return ids[:limit] if limit is not None else ids
|
|
|
|
def batch_update(self, lock_timeout: float = 30.0) -> WriteBatch:
|
|
"""
|
|
Get a batch context manager for bulk index operations.
|
|
|
|
Use this for efficient bulk document updates/deletions. All operations
|
|
within the batch are committed atomically at the end of the context.
|
|
|
|
Args:
|
|
lock_timeout: Seconds to wait for file lock acquisition
|
|
|
|
Returns:
|
|
WriteBatch context manager
|
|
|
|
Raises:
|
|
SearchIndexLockError: If lock cannot be acquired within timeout
|
|
"""
|
|
self._ensure_open()
|
|
return WriteBatch(self, lock_timeout)
|
|
|
|
def rebuild(
|
|
self,
|
|
documents: QuerySet[Document],
|
|
iter_wrapper: IterWrapper[tuple[Document, ViewerGrant]] = identity,
|
|
writer_heap_bytes: int = 512_000_000,
|
|
) -> None:
|
|
"""
|
|
Rebuild the entire search index from scratch.
|
|
|
|
Wipes the existing index and re-indexes all provided documents.
|
|
On failure, restores the previous index state to keep the backend usable.
|
|
|
|
Args:
|
|
documents: QuerySet of Document instances to index
|
|
iter_wrapper: Optional wrapper function for progress tracking
|
|
(e.g., progress bar). Wraps an iterable of
|
|
``(document, (viewer_ids, viewer_group_ids))`` pairs and should yield
|
|
each unchanged, advancing one step per document.
|
|
writer_heap_bytes: Tantivy writer memory budget (split across the
|
|
writer's threads). Larger values buffer more docs in RAM before
|
|
flushing a segment, deferring merge work; they do not avoid it.
|
|
"""
|
|
# Create new index (on-disk or in-memory)
|
|
if self._path is not None:
|
|
wipe_index(self._path)
|
|
new_index = tantivy.Index(build_schema(), path=str(self._path))
|
|
_write_sentinels(self._path)
|
|
else:
|
|
new_index = tantivy.Index(build_schema())
|
|
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
|
|
|
|
|
|
def chunked(iterable, size):
|
|
iterator = iter(iterable)
|
|
while chunk := list(islice(iterator, size)):
|
|
yield chunk
|
|
|
|
|
|
_EMPTY_VIEWER_GRANT: Final[ViewerGrant] = ViewerGrant(
|
|
viewer_ids=[],
|
|
viewer_group_ids=[],
|
|
)
|
|
|
|
|
|
class _DocumentViewerStream(QuerySetStream["Document"]):
|
|
"""Yield document permission data while batch-loading grants.
|
|
|
|
Viewer permissions are fetched in batches (see
|
|
``_bulk_get_viewer_permissions``), but documents are yielded individually so a
|
|
progress bar wrapped around this stream advances per document rather than
|
|
jumping a whole chunk at a time. ``__len__`` (inherited from
|
|
``QuerySetStream``) lets the progress helper still discover the total (it
|
|
inspects ``QuerySet``/``Sized``).
|
|
|
|
The viewer and group ids travel with each document in the yielded pair
|
|
rather than through a separate mutable attribute, so the pairing survives
|
|
regardless of how ``iter_wrapper`` consumes the stream (buffering,
|
|
batching, etc.) — there is no reliance on the caller advancing this
|
|
generator in lock-step.
|
|
"""
|
|
|
|
def __iter__(self) -> Iterator[tuple[Document, ViewerGrant]]:
|
|
# iterator(chunk_size=…) streams from a server-side cursor instead of
|
|
# materialising the whole queryset in memory; since Django 4.1 it still
|
|
# honours prefetch_related, running the prefetches one batch at a time.
|
|
documents = self._queryset.iterator(chunk_size=self._chunk_size)
|
|
for chunk in chunked(documents, self._chunk_size):
|
|
grants_by_pk = _bulk_get_viewer_permissions([doc.pk for doc in chunk])
|
|
for doc in chunk:
|
|
yield doc, grants_by_pk.get(doc.pk, _EMPTY_VIEWER_GRANT)
|
|
|
|
|
|
def _bulk_get_viewer_permissions(
|
|
doc_pks: Sequence[int],
|
|
) -> dict[int, ViewerGrant]:
|
|
"""Fetch direct user and group view grants for a batch of documents, keyed by pk.
|
|
|
|
Group grants remain group IDs in the index so permission checks use the
|
|
requesting user's current memberships. Expanding groups to user IDs here
|
|
would leave stale access behind after a user is removed from a group.
|
|
"""
|
|
from collections import defaultdict
|
|
|
|
from django.contrib.contenttypes.models import ContentType
|
|
from guardian.models import GroupObjectPermission
|
|
from guardian.models import UserObjectPermission
|
|
|
|
from documents.models import Document
|
|
|
|
# get_for_model is cached by Django, so this costs at most one query total.
|
|
ct = ContentType.objects.get_for_model(Document)
|
|
str_pks = [str(pk) for pk in doc_pks]
|
|
|
|
viewer_map: dict[int, set[int]] = defaultdict(set)
|
|
viewer_group_map: dict[int, set[int]] = defaultdict(set)
|
|
|
|
# Fold the permission lookup into the query via a join on codename instead
|
|
# of a separate Permission.objects.get(), which would otherwise run once per
|
|
# chunk during a full reindex.
|
|
user_qs = UserObjectPermission.objects.filter(
|
|
content_type=ct,
|
|
permission__content_type=ct,
|
|
permission__codename="view_document",
|
|
object_pk__in=str_pks,
|
|
).values_list("object_pk", "user_id")
|
|
for object_pk, user_id in user_qs:
|
|
viewer_map[int(object_pk)].add(user_id)
|
|
|
|
group_qs = GroupObjectPermission.objects.filter(
|
|
content_type=ct,
|
|
permission__content_type=ct,
|
|
permission__codename="view_document",
|
|
object_pk__in=str_pks,
|
|
).values_list("object_pk", "group_id")
|
|
for object_pk, group_id in group_qs:
|
|
viewer_group_map[int(object_pk)].add(group_id)
|
|
|
|
return {
|
|
object_pk: ViewerGrant(
|
|
viewer_ids=list(viewer_map.get(object_pk, ())),
|
|
viewer_group_ids=list(viewer_group_map.get(object_pk, ())),
|
|
)
|
|
for object_pk in viewer_map.keys() | viewer_group_map.keys()
|
|
}
|
|
|
|
|
|
# Module-level singleton with proper thread safety
|
|
_backend: TantivyBackend | None = None
|
|
_backend_path: Path | None = None # tracks which INDEX_DIR the singleton uses
|
|
_backend_lock = threading.RLock()
|
|
|
|
|
|
def get_backend() -> TantivyBackend:
|
|
"""
|
|
Get the global backend instance with thread safety.
|
|
|
|
Returns a singleton TantivyBackend instance, automatically reinitializing
|
|
when settings.INDEX_DIR changes. This ensures proper test isolation when
|
|
using pytest-xdist or @override_settings that change the index directory.
|
|
|
|
Returns:
|
|
Thread-safe singleton TantivyBackend instance
|
|
"""
|
|
global _backend, _backend_path
|
|
|
|
current_path: Path = settings.INDEX_DIR
|
|
|
|
# Fast path: backend is initialized and path hasn't changed (no lock needed)
|
|
if _backend is not None and _backend_path == current_path:
|
|
return _backend
|
|
|
|
# Slow path: first call, or INDEX_DIR changed between calls
|
|
with _backend_lock:
|
|
# Double-check after acquiring lock — another thread may have beaten us
|
|
if _backend is not None and _backend_path == current_path:
|
|
return _backend # pragma: no cover
|
|
|
|
if _backend is not None:
|
|
_backend.close()
|
|
|
|
_backend = TantivyBackend(path=current_path)
|
|
_backend.open()
|
|
_backend_path = current_path
|
|
|
|
return _backend
|
|
|
|
|
|
def reset_backend() -> None:
|
|
"""
|
|
Reset the global backend instance with thread safety.
|
|
|
|
Forces creation of a new backend instance on the next get_backend() call.
|
|
Used for test isolation and when switching between different index directories.
|
|
"""
|
|
global _backend, _backend_path
|
|
|
|
with _backend_lock:
|
|
if _backend is not None:
|
|
_backend.close()
|
|
_backend = None
|
|
_backend_path = None
|