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

1from __future__ import annotations 

2 

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 

19 

20import filelock 

21import tantivy 

22from django.conf import settings 

23from django.utils.timezone import get_current_timezone 

24 

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 

41 

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 

47 

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 

52 

53 from documents.models import Document 

54 

55logger = logging.getLogger("paperless.search") 

56 

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 

61 

62T = TypeVar("T") 

63 

64 

65class ViewerGrant(NamedTuple): 

66 """Direct user and group view grants for a single document. 

67 

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 """ 

72 

73 viewer_ids: list[int] 

74 viewer_group_ids: list[int] 

75 

76 

77class SearchMode(StrEnum): 

78 QUERY = "query" 

79 TEXT = "text" 

80 TITLE = "title" 

81 

82 

83def _extract_autocomplete_words(text_sources: list[str]) -> set[str]: 

84 """Extract and normalize words for autocomplete. 

85 

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 

97 

98 

99class SearchHit(TypedDict): 

100 """Type definition for search result hits.""" 

101 

102 id: int 

103 score: float 

104 rank: int 

105 highlights: dict[str, str] 

106 

107 

108class TantivyRelevanceList: 

109 """ 

110 DRF-compatible list wrapper for Tantivy search results. 

111 

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. 

117 

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 """ 

123 

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 

133 

134 def __len__(self) -> int: 

135 return len(self._ordered_ids) 

136 

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 ] 

161 

162 def get_all_ids(self) -> list[int]: 

163 """Return all matching document IDs in display order.""" 

164 return self._ordered_ids 

165 

166 

167class SearchIndexLockError(Exception): 

168 """Raised when the search index file lock cannot be acquired within the timeout.""" 

169 

170 

171class WriteBatch: 

172 """ 

173 Context manager for bulk index operations with file locking. 

174 

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. 

178 

179 Usage: 

180 with backend.batch_update() as batch: 

181 batch.add_or_update(document) 

182 batch.remove(doc_id) 

183 """ 

184 

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 

190 

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 

197 

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) 

222 

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 

243 

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() 

263 

264 def add_or_update(self, document: Document) -> None: 

265 """ 

266 Add or update a document in the batch. 

267 

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. 

271 

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) 

278 

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 ) 

284 

285 def add_or_update_ids(self, ids: Sequence[int]) -> None: 

286 """ 

287 Add or update multiple documents in the batch by primary key. 

288 

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. 

295 

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. 

300 

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 

306 

307 ids = list(ids) 

308 if not ids: 

309 return 

310 

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) 

330 

331 

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. 

339 

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 

346 

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 

351 

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 ) 

371 

372 

373class TantivyBackend: 

374 """ 

375 Tantivy search backend with explicit lifecycle management. 

376 

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. 

380 

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 """ 

384 

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 } 

397 

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 ) 

412 

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 

417 

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 

422 

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 

427 

428 def open(self) -> None: 

429 """ 

430 Open or rebuild the index as needed. 

431 

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 

441 

442 def close(self) -> None: 

443 """ 

444 Close the index and release resources. 

445 

446 Safe to call multiple times - subsequent calls are no-ops. 

447 """ 

448 self._raw_index = None 

449 self._raw_schema = None 

450 

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 

455 

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) 

470 

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 

486 

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 ) 

495 

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. 

503 

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 

511 

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) 

517 

518 doc = tantivy.Document() 

519 

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) 

534 

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 ) 

541 

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) 

549 

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) 

557 

558 # Storage path 

559 if document.storage_path: 

560 doc.add_text( 

561 "storage_path", 

562 normalize_search_text(document.storage_path.name), 

563 ) 

564 

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) 

573 

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)) 

596 

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 ) 

612 

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 ) 

622 

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) 

633 

634 if document.archive_serial_number is not None: 

635 doc.add_unsigned("asn", document.archive_serial_number) 

636 

637 if document.page_count is not None: 

638 doc.add_unsigned("page_count", document.page_count) 

639 

640 doc.add_unsigned("num_notes", num_notes) 

641 

642 # Owner 

643 if document.owner_id: 

644 doc.add_unsigned("owner_id", document.owner_id) 

645 

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) 

671 

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) 

679 

680 for word in sorted(_extract_autocomplete_words(text_sources)): 

681 doc.add_text("autocomplete_word", word) 

682 

683 return doc 

684 

685 def add_or_update(self, document: Document) -> None: 

686 """ 

687 Add or update a single document with file locking. 

688 

689 Convenience method for single-document updates. For bulk operations, 

690 use batch_update() context manager for better performance. 

691 

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. 

695 

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 

711 

712 index_document.apply_async(args=[document.pk], countdown=60) 

713 

714 def remove(self, doc_id: int) -> None: 

715 """ 

716 Remove a single document from the index with file locking. 

717 

718 Convenience method for single-document removal. For bulk operations, 

719 use batch_update() context manager for better performance. 

720 

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. 

724 

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 

740 

741 remove_document_from_index.apply_async(args=[doc_id], countdown=60) 

742 

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. 

753 

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. 

758 

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) 

765 

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 [] 

771 

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 ) 

790 

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 

806 

807 searcher = self._index.searcher() 

808 

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)) 

827 

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 } 

834 

835 snippet_generator = None 

836 notes_snippet_generator = None 

837 hits: list[SearchHit] = [] 

838 

839 for rank, doc_id in enumerate(doc_ids, start=rank_start): 

840 if doc_id not in addr_by_id: 

841 continue 

842 

843 score, doc_address = addr_by_id[doc_id] 

844 actual_doc = searcher.doc(doc_address) 

845 doc_dict = actual_doc.to_dict() 

846 

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 ) 

856 

857 content_html = snippet_generator.snippet_from_doc(actual_doc).to_html() 

858 if content_html: 

859 highlights["content"] = content_html 

860 

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 

876 

877 except Exception: # pragma: no cover 

878 logger.debug("Failed to generate highlights for doc %s", doc_id) 

879 

880 hits.append( 

881 SearchHit( 

882 id=doc_id, 

883 score=score, 

884 rank=rank, 

885 highlights=highlights, 

886 ), 

887 ) 

888 

889 return hits 

890 

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. 

903 

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. 

907 

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) 

915 

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) 

922 

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 [] 

927 

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] 

940 

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] 

945 

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] 

949 

950 return cast( 

951 "list[int]", 

952 searcher.fast_field_values("id", [doc_addr for doc_addr, *_ in all_hits]), 

953 ) 

954 

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. 

963 

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. 

967 

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. 

971 

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) 

976 

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 [] 

984 

985 searcher = self._index.searcher() 

986 

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) 

992 

993 matches = searcher.terms_with_prefix( 

994 "autocomplete_word", 

995 normalized_term, 

996 permission_query, 

997 limit, 

998 ) 

999 

1000 return [x[0] for x in matches] 

1001 

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. 

1011 

1012 Lightweight companion to more_like_this(). The original document is 

1013 excluded from results. 

1014 

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) 

1019 

1020 Returns: 

1021 List of similar document IDs (excluding the original) 

1022 """ 

1023 self._ensure_open() 

1024 searcher = self._index.searcher() 

1025 

1026 id_query = tantivy.Query.term_query(self._schema, "id", doc_id) 

1027 results = searcher.search(id_query, limit=1) 

1028 

1029 if not results.hits: 

1030 return [] 

1031 

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 ) 

1043 

1044 final_query = self._apply_permission_filter(mlt_query, user) 

1045 

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 [] 

1065 

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 

1070 

1071 def batch_update(self, lock_timeout: float = 30.0) -> WriteBatch: 

1072 """ 

1073 Get a batch context manager for bulk index operations. 

1074 

1075 Use this for efficient bulk document updates/deletions. All operations 

1076 within the batch are committed atomically at the end of the context. 

1077 

1078 Args: 

1079 lock_timeout: Seconds to wait for file lock acquisition 

1080 

1081 Returns: 

1082 WriteBatch context manager 

1083 

1084 Raises: 

1085 SearchIndexLockError: If lock cannot be acquired within timeout 

1086 """ 

1087 self._ensure_open() 

1088 return WriteBatch(self, lock_timeout) 

1089 

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. 

1098 

1099 Wipes the existing index and re-indexes all provided documents. 

1100 On failure, restores the previous index state to keep the backend usable. 

1101 

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) 

1116 

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 

1146 

1147 

1148def chunked(iterable, size): 

1149 iterator = iter(iterable) 

1150 while chunk := list(islice(iterator, size)): 

1151 yield chunk 

1152 

1153 

1154_EMPTY_VIEWER_GRANT: Final[ViewerGrant] = ViewerGrant( 

1155 viewer_ids=[], 

1156 viewer_group_ids=[], 

1157) 

1158 

1159 

1160class _DocumentViewerStream(QuerySetStream["Document"]): 

1161 """Yield document permission data while batch-loading grants. 

1162 

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``). 

1169 

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 """ 

1176 

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) 

1186 

1187 

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. 

1192 

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 

1198 

1199 from django.contrib.contenttypes.models import ContentType 

1200 from guardian.models import GroupObjectPermission 

1201 from guardian.models import UserObjectPermission 

1202 

1203 from documents.models import Document 

1204 

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] 

1208 

1209 viewer_map: dict[int, set[int]] = defaultdict(set) 

1210 viewer_group_map: dict[int, set[int]] = defaultdict(set) 

1211 

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) 

1223 

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) 

1232 

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 } 

1240 

1241 

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() 

1246 

1247 

1248def get_backend() -> TantivyBackend: 

1249 """ 

1250 Get the global backend instance with thread safety. 

1251 

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. 

1255 

1256 Returns: 

1257 Thread-safe singleton TantivyBackend instance 

1258 """ 

1259 global _backend, _backend_path 

1260 

1261 current_path: Path = settings.INDEX_DIR 

1262 

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 

1266 

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 

1272 

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() 

1275 

1276 _backend = TantivyBackend(path=current_path) 

1277 _backend.open() 

1278 _backend_path = current_path 

1279 

1280 return _backend 

1281 

1282 

1283def reset_backend() -> None: 

1284 """ 

1285 Reset the global backend instance with thread safety. 

1286 

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 

1291 

1292 with _backend_lock: 

1293 if _backend is not None: 

1294 _backend.close() 

1295 _backend = None 

1296 _backend_path = None