Coverage for paperless_ai/indexing.py: 18%

279 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 09:07 +0000

1import enum 

2import logging 

3from collections.abc import Iterable 

4from contextlib import contextmanager 

5from datetime import timedelta 

6from typing import TYPE_CHECKING 

7 

8from django.conf import settings 

9from django.core.exceptions import ObjectDoesNotExist 

10from django.utils import timezone 

11from filelock import FileLock 

12from filelock import ReadWriteLock 

13from filelock import Timeout 

14 

15from documents.models import Document 

16from documents.models import PaperlessTask 

17from documents.utils import IterWrapper 

18from documents.utils import QuerySetStream 

19from documents.utils import identity 

20from paperless.config import AIConfig 

21from paperless_ai.db import db_connection_released 

22from paperless_ai.embedding import build_llm_index_text 

23from paperless_ai.embedding import get_configured_model_name 

24from paperless_ai.embedding import get_embedding_model 

25 

26if TYPE_CHECKING: 26 ↛ 27line 26 didn't jump to line 27 because the condition on line 26 was never true

27 from llama_index.core.schema import BaseNode 

28 from llama_index.core.schema import NodeWithScore 

29 

30 from paperless_ai.vector_store import PaperlessSqliteVecVectorStore 

31 

32 

33logger = logging.getLogger("paperless_ai.indexing") 

34 

35RAG_NUM_OUTPUT = 512 

36RAG_CHUNK_OVERLAP = 200 

37 

38# update_llm_index(): row count per .iterator() batch when streaming 

39# documents for a rebuild/update via QuerySetStream, matching 

40# _DocumentViewerStream's chunk size in documents/search/_backend.py. 

41_INDEX_STREAM_CHUNK_SIZE = 1000 

42 

43 

44def queue_llm_index_update_if_needed(*, rebuild: bool, reason: str) -> bool: 

45 # NOTE: The check-then-enqueue sequence below is non-atomic (TOCTOU): two 

46 # concurrent workers can both observe no running task and both enqueue a 

47 # full rebuild. This is wasteful but not data-corrupting — update_llm_index 

48 # is itself protected by settings.LLM_INDEX_LOCK, so only one rebuild runs at a 

49 # time and the second one is serialised after the first completes. 

50 from documents.tasks import llmindex_index 

51 

52 has_running = PaperlessTask.objects.filter( 

53 task_type=PaperlessTask.TaskType.LLM_INDEX, 

54 status__in=[PaperlessTask.Status.PENDING, PaperlessTask.Status.STARTED], 

55 ).exists() 

56 has_recent = PaperlessTask.objects.filter( 

57 task_type=PaperlessTask.TaskType.LLM_INDEX, 

58 date_created__gte=(timezone.now() - timedelta(minutes=5)), 

59 ).exists() 

60 if has_running or has_recent: 

61 return False 

62 

63 llmindex_index.apply_async( 

64 kwargs={"rebuild": rebuild}, 

65 headers={"trigger_source": PaperlessTask.TriggerSource.SYSTEM}, 

66 ) 

67 logger.warning( 

68 "Queued LLM index update%s: %s", 

69 " (rebuild)" if rebuild else "", 

70 reason, 

71 ) 

72 return True 

73 

74 

75def get_vector_store() -> "PaperlessSqliteVecVectorStore": 

76 from paperless_ai.vector_store import PaperlessSqliteVecVectorStore 

77 

78 settings.LLM_INDEX_DIR.mkdir(parents=True, exist_ok=True) 

79 return PaperlessSqliteVecVectorStore( 

80 uri=str(settings.LLM_INDEX_DIR), 

81 ) 

82 

83 

84# --- LLM index locking --------------------------------------------------- 

85# 

86# Two locks guard the index; they answer different questions and are NOT 

87# interchangeable: 

88# 

89# * settings.LLM_INDEX_LOCK (FileLock, exclusive) - serializes WRITERS against 

90# each other, so only one rebuild/upsert/delete/compaction runs at a time. 

91# Taken by write_store(). Readers never take it, so it never blocks reads. 

92# 

93# * settings.LLM_INDEX_RWLOCK (ReadWriteLock) - coordinates readers against the 

94# compaction/migration file swap. read_store() takes it SHARED (readers run 

95# concurrently); _exclude_readers() takes it EXCLUSIVE, only for the swap, so 

96# the database file is never replaced while a reader connection is open (that 

97# would alias the old WAL onto the new file and corrupt it). 

98# 

99# | vs another writer | vs a reader 

100# -----------------+-------------------+---------------------------- 

101# normal write | LLM_INDEX_LOCK | nothing (WAL gives MVCC) 

102# compaction/swap | LLM_INDEX_LOCK | LLM_INDEX_RWLOCK (exclusive) 

103# reader | nothing (WAL) | LLM_INDEX_RWLOCK (shared) 

104# 

105# They can't be merged into one ReadWriteLock: a normal write must exclude other 

106# writers WITHOUT blocking readers (WAL already gives reader/writer concurrency), 

107# and ReadWriteLock has no "exclusive vs writers, shared vs readers" mode. Only 

108# the swap needs to exclude readers. 

109def _index_rwlock() -> ReadWriteLock: 

110 """Return a fresh read/write lock instance for the index swap. 

111 

112 ``is_singleton=False`` so reads and the swap always coordinate through 

113 SQLite (the actual cross-process case) rather than hitting the in-process 

114 reentrant-upgrade guard; callers must ``close()`` it (the context managers 

115 below do). 

116 """ 

117 settings.LLM_INDEX_DIR.mkdir(parents=True, exist_ok=True) 

118 return ReadWriteLock(str(settings.LLM_INDEX_RWLOCK), is_singleton=False) 

119 

120 

121@contextmanager 

122def read_store(): 

123 """Acquire the shared read lock and yield the vector store for a read. 

124 

125 The shared lock is held for the whole lifetime of the connection (and 

126 closed on exit) so the compaction/migration swap, which takes the exclusive 

127 lock, never runs while this connection is open. Concurrent readers do not 

128 block each other; only the swap does. 

129 """ 

130 lock = _index_rwlock() 

131 try: 

132 with lock.read_lock(), get_vector_store() as store: 

133 yield store 

134 finally: 

135 lock.close() 

136 

137 

138@contextmanager 

139def _exclude_readers(): 

140 """Acquire exclusive index access, blocking until readers have drained. 

141 

142 The exclusive counterpart to ``read_store()``: a compaction or migration 

143 must not run while any reader connection is open. Raises 

144 :class:`filelock.Timeout` if active readers do not drain within 

145 ``LLM_INDEX_COMPACTION_LOCK_TIMEOUT``; callers skip the operation on timeout. 

146 """ 

147 lock = _index_rwlock() 

148 try: 

149 with lock.write_lock(timeout=settings.LLM_INDEX_COMPACTION_LOCK_TIMEOUT): 

150 yield 

151 finally: 

152 lock.close() 

153 

154 

155def _with_exclusive_access(operation: str, fn): 

156 """Run ``fn()`` with exclusive index access (see ``_exclude_readers()``), 

157 for compaction/migration file swaps that must not run while readers are 

158 active. Returns ``fn()``'s result, or None (after logging) if active 

159 readers do not drain within ``LLM_INDEX_COMPACTION_LOCK_TIMEOUT`` -- 

160 callers skip the operation this run; it retries next time. 

161 """ 

162 try: 

163 with _exclude_readers(): 

164 return fn() 

165 except Timeout: 

166 logger.info( 

167 "Skipping LLM index %s: index readers are active; will retry next run.", 

168 operation, 

169 ) 

170 return None 

171 

172 

173@contextmanager 

174def write_store(embed_model_name: str | None = None): 

175 """Acquire the write lock and yield the vector store. 

176 

177 All mutating operations (upsert, delete, rebuild, compact) must go through 

178 this context manager to serialise concurrent Celery writers. 

179 Read paths use ``read_store()`` so they hold the shared read lock. 

180 

181 Pass ``embed_model_name`` whenever the operation may create the table so 

182 the model name is recorded in the schema metadata for future mismatch checks. 

183 """ 

184 from paperless_ai.vector_store import PaperlessSqliteVecVectorStore 

185 

186 settings.LLM_INDEX_DIR.mkdir(parents=True, exist_ok=True) 

187 with ( 

188 FileLock(settings.LLM_INDEX_LOCK), 

189 PaperlessSqliteVecVectorStore( 

190 uri=str(settings.LLM_INDEX_DIR), 

191 embed_model_name=embed_model_name, 

192 ) as store, 

193 ): 

194 yield store 

195 

196 

197class MigrationCheckResult(enum.Enum): 

198 """Outcome of _check_and_run_migrations(). 

199 

200 CURRENT: no migration was pending, or a pending structural migration 

201 was applied successfully - safe to write. 

202 

203 REEMBED_REQUIRED: a pending migration needs fresh embeddings, which is 

204 never triggered automatically - the caller must force a rebuild. 

205 

206 DEFERRED: a migration was pending but could not run because active 

207 index readers did not drain within LLM_INDEX_COMPACTION_LOCK_TIMEOUT -- 

208 the store is still on its old schema. Callers must NOT proceed to 

209 write: collapsing this into the same falsy value as CURRENT (as a 

210 plain bool return once did) would let a write proceed against an 

211 unmigrated schema. 

212 """ 

213 

214 CURRENT = "current" 

215 REEMBED_REQUIRED = "reembed_required" 

216 DEFERRED = "deferred" 

217 

218 

219def _check_and_run_migrations( 

220 store: "PaperlessSqliteVecVectorStore", 

221) -> MigrationCheckResult: 

222 """Run any pending structural migrations, reporting the outcome as a 

223 tri-state result. Safe to call before any write, including 

224 delete()/upsert_document(): has_pending_migration() (see its docstring) 

225 keeps this a no-op, with no exclusive access taken, once the store is 

226 current. 

227 """ 

228 if not store.has_pending_migration(): 

229 return MigrationCheckResult.CURRENT 

230 result = _with_exclusive_access( 

231 "migration check", 

232 store.check_and_run_migrations, 

233 ) 

234 if result is None: 

235 return MigrationCheckResult.DEFERRED 

236 return ( 

237 MigrationCheckResult.REEMBED_REQUIRED 

238 if result 

239 else MigrationCheckResult.CURRENT 

240 ) 

241 

242 

243def _safe_related_name(document: Document, field: str) -> str | None: 

244 """ 

245 Returns the ``name`` of a related object (correspondent, document_type, 

246 storage_path), or None if the FK is unset or points at a row that has 

247 since been deleted (e.g. concurrently with this call). 

248 """ 

249 try: 

250 related = getattr(document, field) 

251 except ObjectDoesNotExist: 

252 return None 

253 return related.name if related else None 

254 

255 

256def build_document_node( 

257 document: Document, 

258 *, 

259 chunk_size: int | None = None, 

260) -> list["BaseNode"]: 

261 """ 

262 Given a Document, returns parsed Nodes ready for indexing. 

263 """ 

264 text = build_llm_index_text(document) 

265 metadata = { 

266 "document_id": str(document.id), 

267 "title": document.title, 

268 "tags": [t.name for t in document.tags.all()], 

269 "correspondent": _safe_related_name(document, "correspondent"), 

270 "document_type": _safe_related_name(document, "document_type"), 

271 "filename": document.filename, 

272 "storage_path": _safe_related_name(document, "storage_path"), 

273 "archive_serial_number": document.archive_serial_number, 

274 "created": document.created.isoformat() if document.created else None, 

275 "added": document.added.isoformat() if document.added else None, 

276 "modified": document.modified.isoformat(), 

277 } 

278 from llama_index.core import Document as LlamaDocument 

279 from llama_index.core.node_parser import SimpleNodeParser 

280 

281 # Exclude all metadata keys from the embedding text — build_llm_index_text 

282 # already encodes this info in the body, so prepending it again would double 

283 # the token count and exceed embedding models with small context windows 

284 # (e.g. nomic-embed-text via Ollama defaults to num_ctx=2048). 

285 doc = LlamaDocument( 

286 id_=str(document.id), 

287 text=text, 

288 metadata=metadata, 

289 excluded_embed_metadata_keys=list(metadata.keys()), 

290 excluded_llm_metadata_keys=["document_id"], 

291 ) 

292 chunk_size = chunk_size or get_rag_chunk_size() 

293 parser = SimpleNodeParser( 

294 chunk_size=chunk_size, 

295 chunk_overlap=get_rag_chunk_overlap(chunk_size), 

296 ) 

297 return parser.get_nodes_from_documents([doc]) 

298 

299 

300def load_or_build_index(config: AIConfig, store: "PaperlessSqliteVecVectorStore"): 

301 """Return a VectorStoreIndex backed by ``store``. 

302 

303 ``store`` is supplied by the caller's ``read_store()`` context so the shared 

304 read lock and the connection stay alive for the whole retrieval. 

305 """ 

306 import llama_index.core.settings as llama_settings 

307 from llama_index.core import VectorStoreIndex 

308 

309 embed_model = get_embedding_model(config) 

310 llama_settings.Settings.embed_model = embed_model 

311 return VectorStoreIndex.from_vector_store( 

312 vector_store=store, 

313 embed_model=embed_model, 

314 ) 

315 

316 

317def llm_index_exists() -> bool: 

318 """True when the index table exists on disk.""" 

319 with read_store() as store: 

320 return store.table_exists() 

321 

322 

323def get_rag_chunk_size() -> int: 

324 return AIConfig().llm_embedding_chunk_size 

325 

326 

327def get_rag_chunk_overlap(chunk_size: int | None = None) -> int: 

328 chunk_size = chunk_size or get_rag_chunk_size() 

329 return min(RAG_CHUNK_OVERLAP, chunk_size - 1) 

330 

331 

332def get_rag_prompt_helper( 

333 *, 

334 chunk_size: int | None = None, 

335 context_size: int | None = None, 

336): 

337 from llama_index.core.indices.prompt_helper import PromptHelper 

338 

339 if chunk_size is None or context_size is None: 

340 config = AIConfig() 

341 chunk_size = chunk_size or config.llm_embedding_chunk_size 

342 context_size = context_size or config.llm_context_size 

343 

344 return PromptHelper( 

345 context_window=context_size, 

346 num_output=RAG_NUM_OUTPUT, 

347 chunk_overlap_ratio=0.1, 

348 chunk_size_limit=chunk_size, 

349 ) 

350 

351 

352def _embed_nodes(nodes: list["BaseNode"], embed_model) -> None: 

353 """Embed ``nodes`` in place using ``embed_model``.""" 

354 from llama_index.core.schema import MetadataMode 

355 

356 texts = [n.get_content(metadata_mode=MetadataMode.EMBED) for n in nodes] 

357 for node, emb in zip( 

358 nodes, 

359 embed_model.get_text_embedding_batch(texts), 

360 strict=True, 

361 ): 

362 node.embedding = emb 

363 

364 

365def document_id_filters(doc_ids): 

366 """Return a MetadataFilters IN filter scoped to ``doc_ids``.""" 

367 from llama_index.core.vector_stores.types import FilterOperator 

368 from llama_index.core.vector_stores.types import MetadataFilter 

369 from llama_index.core.vector_stores.types import MetadataFilters 

370 

371 return MetadataFilters( 

372 filters=[ 

373 MetadataFilter( 

374 key="document_id", 

375 operator=FilterOperator.IN, 

376 value=list(doc_ids), 

377 ), 

378 ], 

379 ) 

380 

381 

382def _exclude_document_id_filter(document_id: int | str): 

383 """Return a MetadataFilters NE filter excluding ``document_id``.""" 

384 from llama_index.core.vector_stores.types import FilterOperator 

385 from llama_index.core.vector_stores.types import MetadataFilter 

386 from llama_index.core.vector_stores.types import MetadataFilters 

387 

388 return MetadataFilters( 

389 filters=[ 

390 MetadataFilter( 

391 key="document_id", 

392 operator=FilterOperator.NE, 

393 value=str(document_id), 

394 ), 

395 ], 

396 ) 

397 

398 

399def exclude_document_ids_filter(doc_ids): 

400 """Return a MetadataFilters NIN filter excluding every id in ``doc_ids``.""" 

401 from llama_index.core.vector_stores.types import FilterOperator 

402 from llama_index.core.vector_stores.types import MetadataFilter 

403 from llama_index.core.vector_stores.types import MetadataFilters 

404 

405 return MetadataFilters( 

406 filters=[ 

407 MetadataFilter( 

408 key="document_id", 

409 operator=FilterOperator.NIN, 

410 value=list(doc_ids), 

411 ), 

412 ], 

413 ) 

414 

415 

416def update_llm_index( 

417 *, 

418 iter_wrapper: IterWrapper[Document] = identity, 

419 rebuild=False, 

420 document_ids: Iterable[int] | None = None, 

421) -> str: 

422 """Rebuild or incrementally update the LLM index. 

423 

424 ``document_ids``, when given, scopes an incremental update to just those 

425 documents instead of scanning the whole library - callers that already 

426 know which documents changed (e.g. a bulk edit) should pass this to avoid 

427 an O(library size) scan per call. Ignored whenever a rebuild actually 

428 happens, since a rebuild always covers the whole library regardless. 

429 """ 

430 with write_store() as store: 

431 migration_result = _check_and_run_migrations(store) 

432 if migration_result is MigrationCheckResult.REEMBED_REQUIRED: 

433 logger.warning( 

434 "LLM index migration requires re-embedding; forcing rebuild.", 

435 ) 

436 rebuild = True 

437 elif migration_result is MigrationCheckResult.DEFERRED: 

438 logger.info( 

439 "Skipping LLM index update: migration check deferred while " 

440 "index readers are active; will retry next run.", 

441 ) 

442 return ( 

443 "Skipping LLM index update: migration check deferred; " 

444 "will retry next run." 

445 ) 

446 documents = Document.objects.select_related( 

447 "correspondent", 

448 "document_type", 

449 "storage_path", 

450 ).prefetch_related("tags", "notes", "custom_fields__field") 

451 no_documents = not documents.exists() 

452 

453 # Fast exit before touching config: nothing to index and no existing index. 

454 if no_documents and not rebuild and not llm_index_exists(): 

455 logger.warning("No documents found to index.") 

456 return "No documents found to index." 

457 

458 config = AIConfig() 

459 model_name = get_configured_model_name(config) 

460 

461 if not rebuild: 

462 with read_store() as store: 

463 if store.table_exists() and store.config_mismatch(model_name): 

464 logger.warning("Embedding model changed; forcing LLM index rebuild.") 

465 rebuild = True 

466 

467 if no_documents: 

468 logger.warning("No documents found to index.") 

469 

470 chunk_size = config.llm_embedding_chunk_size 

471 embed_model = get_embedding_model(config) 

472 

473 with write_store(embed_model_name=model_name) as store: 

474 if rebuild or not store.table_exists(): 

475 logger.info("Rebuilding LLM index.") 

476 store.drop_table() 

477 for document in iter_wrapper( 

478 QuerySetStream(documents, chunk_size=_INDEX_STREAM_CHUNK_SIZE), 

479 ): 

480 nodes = build_document_node(document, chunk_size=chunk_size) 

481 _embed_nodes(nodes, embed_model) 

482 store.add(nodes) 

483 msg = "LLM index rebuilt successfully." 

484 else: 

485 scoped_documents = ( 

486 documents.filter(id__in=document_ids) 

487 if document_ids is not None 

488 else documents 

489 ) 

490 existing = store.get_modified_times() 

491 changed = 0 

492 for document in iter_wrapper( 

493 QuerySetStream(scoped_documents, chunk_size=_INDEX_STREAM_CHUNK_SIZE), 

494 ): 

495 doc_id = str(document.id) 

496 if existing.get(doc_id) == document.modified.isoformat(): 

497 continue 

498 nodes = build_document_node(document, chunk_size=chunk_size) 

499 _embed_nodes(nodes, embed_model) 

500 store.upsert_document(doc_id, nodes) 

501 changed += 1 

502 msg = ( 

503 "LLM index updated successfully." 

504 if changed 

505 else "No changes detected in LLM index." 

506 ) 

507 

508 _with_exclusive_access("compaction", store.compact) 

509 return msg 

510 

511 

512def llm_index_add_or_update_document(document: Document): 

513 """Add or atomically replace a document's chunks in the index.""" 

514 config = AIConfig() 

515 new_nodes = build_document_node( 

516 document, 

517 chunk_size=config.llm_embedding_chunk_size, 

518 ) 

519 if new_nodes: 

520 _embed_nodes(new_nodes, get_embedding_model(config)) 

521 

522 with write_store(embed_model_name=get_configured_model_name(config)) as store: 

523 migration_result = _check_and_run_migrations(store) 

524 if migration_result is MigrationCheckResult.REEMBED_REQUIRED: 

525 logger.warning( 

526 "Skipping incremental LLM index update for document %s: the " 

527 "index requires re-embedding first. Run 'document_llmindex " 

528 "rebuild' to resolve.", 

529 document.id, 

530 ) 

531 return 

532 if migration_result is MigrationCheckResult.DEFERRED: 

533 logger.info( 

534 "Skipping incremental LLM index update for document %s: " 

535 "migration check deferred while index readers are active; " 

536 "will retry on the next write.", 

537 document.id, 

538 ) 

539 return 

540 store.upsert_document(str(document.id), new_nodes) 

541 

542 

543def llm_index_migrate() -> None: 

544 """Apply any pending LLM index schema migrations, with no reindex. 

545 

546 Intended to run unconditionally on every startup (see the 

547 init-llmindex-migrate container step and the bare-metal upgrade docs): 

548 has_pending_migration() short-circuits to a metadata-only read once the 

549 store is current, so a healthy install pays almost nothing here. Only 

550 ever applies structural migrations - a pending re-embed migration is 

551 left for the explicit, deliberate rebuild path (``document_llmindex 

552 update``/``rebuild``) to resolve, since re-embedding can be slow and, 

553 for a metered embedding backend, cost money. 

554 """ 

555 if not AIConfig().llm_index_enabled: 

556 return 

557 with write_store() as store: 

558 migration_result = _check_and_run_migrations(store) 

559 if migration_result is MigrationCheckResult.REEMBED_REQUIRED: 

560 logger.warning( 

561 "LLM index requires re-embedding, which this automatic migration " 

562 "check will not do on its own - it can be slow and, for a " 

563 "metered embedding backend, cost money. Run " 

564 "'document_llmindex rebuild' manually when ready.", 

565 ) 

566 elif migration_result is MigrationCheckResult.DEFERRED: 

567 logger.info( 

568 "LLM index migration check deferred while index readers are " 

569 "active; will retry next run.", 

570 ) 

571 

572 

573def llm_index_compact() -> None: 

574 """Compact the index immediately, rebuilding the table to reclaim space.""" 

575 with write_store() as store: 

576 _with_exclusive_access("compaction", lambda: store.compact(force=True)) 

577 

578 

579def llm_index_remove_document(document: Document): 

580 """Remove a document's chunks from the LLM index.""" 

581 with write_store() as store: 

582 migration_result = _check_and_run_migrations(store) 

583 if migration_result is MigrationCheckResult.REEMBED_REQUIRED: 

584 logger.warning( 

585 "Skipping removal of document %s from the LLM index: the " 

586 "index requires re-embedding first. Run 'document_llmindex " 

587 "rebuild' to resolve.", 

588 document.id, 

589 ) 

590 return 

591 if migration_result is MigrationCheckResult.DEFERRED: 

592 logger.info( 

593 "Skipping removal of document %s from the LLM index: " 

594 "migration check deferred while index readers are active; " 

595 "will retry on the next write.", 

596 document.id, 

597 ) 

598 return 

599 store.delete(str(document.id)) 

600 

601 

602def truncate_content( 

603 content: str, 

604 *, 

605 chunk_size: int | None = None, 

606 context_size: int | None = None, 

607) -> str: 

608 from llama_index.core.prompts import PromptTemplate 

609 from llama_index.core.text_splitter import TokenTextSplitter 

610 

611 if chunk_size is None or context_size is None: 

612 config = AIConfig() 

613 chunk_size = chunk_size or config.llm_embedding_chunk_size 

614 context_size = context_size or config.llm_context_size 

615 prompt_helper = get_rag_prompt_helper( 

616 chunk_size=chunk_size, 

617 context_size=context_size, 

618 ) 

619 splitter = TokenTextSplitter( 

620 separator=" ", 

621 chunk_size=chunk_size, 

622 chunk_overlap=get_rag_chunk_overlap(chunk_size), 

623 ) 

624 content_chunks = splitter.split_text(content) 

625 truncated_chunks = prompt_helper.truncate( 

626 prompt=PromptTemplate(template="{content}"), 

627 text_chunks=content_chunks, 

628 padding=5, 

629 ) 

630 return " ".join(truncated_chunks) 

631 

632 

633def truncate_embedding_query(content: str, *, chunk_size: int) -> str: 

634 from llama_index.core.text_splitter import TokenTextSplitter 

635 

636 splitter = TokenTextSplitter( 

637 separator=" ", 

638 chunk_size=chunk_size, 

639 chunk_overlap=0, 

640 ) 

641 content_chunks = splitter.split_text(content) 

642 return content_chunks[0] if content_chunks else "" 

643 

644 

645def normalize_document_ids(document_ids: Iterable[int | str] | None) -> set[str] | None: 

646 if document_ids is None: 

647 return None 

648 return {str(document_id) for document_id in document_ids} 

649 

650 

651def retrieve_similar_nodes( 

652 document: Document, 

653 top_k: int = 5, 

654 document_ids: Iterable[int | str] | None = None, 

655) -> list["NodeWithScore"]: 

656 """Run the vector-store retrieval once and return the raw scored nodes, 

657 permission-filtered by document_ids and with the source document excluded. 

658 Callers derive both RAG text context and taxonomy candidates from this 

659 single retrieval instead of querying the vector store twice per request. 

660 """ 

661 allowed_document_ids = normalize_document_ids(document_ids) 

662 if allowed_document_ids is not None and not allowed_document_ids: 

663 return [] 

664 

665 if not llm_index_exists(): 

666 queue_llm_index_update_if_needed( 

667 rebuild=False, 

668 reason="LLM index not found for similarity query.", 

669 ) 

670 return [] 

671 

672 config = AIConfig() 

673 

674 from llama_index.core.retrievers import VectorIndexRetriever 

675 from llama_index.core.vector_stores.types import FilterCondition 

676 from llama_index.core.vector_stores.types import MetadataFilters 

677 

678 filter_parts = [] 

679 if allowed_document_ids is not None: 

680 filter_parts.extend(document_id_filters(allowed_document_ids).filters) 

681 if document.pk is not None: 

682 filter_parts.extend(_exclude_document_id_filter(document.pk).filters) 

683 

684 filters = ( 

685 MetadataFilters(filters=filter_parts, condition=FilterCondition.AND) 

686 if filter_parts 

687 else None 

688 ) 

689 

690 query_text = truncate_embedding_query( 

691 (document.title or "") + "\n" + (document.content or ""), 

692 chunk_size=config.llm_embedding_chunk_size, 

693 ) 

694 # Hold the shared read lock for the whole retrieval so the connection is 

695 # never open across a compaction swap. The retrieve() call generates a 

696 # query embedding (a slow external request) and searches the vector store; 

697 # no Django ORM access happens during it, so release the pooled DB 

698 # connection for its duration. See #12976. 

699 with read_store() as store: 

700 index = load_or_build_index(config, store) 

701 retriever = VectorIndexRetriever( 

702 index=index, 

703 similarity_top_k=top_k, 

704 filters=filters, 

705 ) 

706 with db_connection_released(): 

707 results = retriever.retrieve(query_text) 

708 

709 if allowed_document_ids is None: 

710 return results 

711 

712 filtered = [] 

713 for node in results: 

714 document_id = node.metadata.get("document_id") 

715 if document_id is None: # pragma: no cover 

716 # Every node the indexing pipeline builds always sets 

717 # document_id; this guards a malformed/partial vec0 row that 

718 # shouldn't occur given the current schema. 

719 continue 

720 if str(document_id) not in allowed_document_ids: 

721 continue 

722 filtered.append(node) 

723 return filtered