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
« 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
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
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
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
30 from paperless_ai.vector_store import PaperlessSqliteVecVectorStore
33logger = logging.getLogger("paperless_ai.indexing")
35RAG_NUM_OUTPUT = 512
36RAG_CHUNK_OVERLAP = 200
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
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
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
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
75def get_vector_store() -> "PaperlessSqliteVecVectorStore":
76 from paperless_ai.vector_store import PaperlessSqliteVecVectorStore
78 settings.LLM_INDEX_DIR.mkdir(parents=True, exist_ok=True)
79 return PaperlessSqliteVecVectorStore(
80 uri=str(settings.LLM_INDEX_DIR),
81 )
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.
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)
121@contextmanager
122def read_store():
123 """Acquire the shared read lock and yield the vector store for a read.
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()
138@contextmanager
139def _exclude_readers():
140 """Acquire exclusive index access, blocking until readers have drained.
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()
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
173@contextmanager
174def write_store(embed_model_name: str | None = None):
175 """Acquire the write lock and yield the vector store.
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.
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
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
197class MigrationCheckResult(enum.Enum):
198 """Outcome of _check_and_run_migrations().
200 CURRENT: no migration was pending, or a pending structural migration
201 was applied successfully - safe to write.
203 REEMBED_REQUIRED: a pending migration needs fresh embeddings, which is
204 never triggered automatically - the caller must force a rebuild.
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 """
214 CURRENT = "current"
215 REEMBED_REQUIRED = "reembed_required"
216 DEFERRED = "deferred"
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 )
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
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
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])
300def load_or_build_index(config: AIConfig, store: "PaperlessSqliteVecVectorStore"):
301 """Return a VectorStoreIndex backed by ``store``.
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
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 )
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()
323def get_rag_chunk_size() -> int:
324 return AIConfig().llm_embedding_chunk_size
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)
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
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
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 )
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
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
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
371 return MetadataFilters(
372 filters=[
373 MetadataFilter(
374 key="document_id",
375 operator=FilterOperator.IN,
376 value=list(doc_ids),
377 ),
378 ],
379 )
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
388 return MetadataFilters(
389 filters=[
390 MetadataFilter(
391 key="document_id",
392 operator=FilterOperator.NE,
393 value=str(document_id),
394 ),
395 ],
396 )
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
405 return MetadataFilters(
406 filters=[
407 MetadataFilter(
408 key="document_id",
409 operator=FilterOperator.NIN,
410 value=list(doc_ids),
411 ),
412 ],
413 )
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.
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()
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."
458 config = AIConfig()
459 model_name = get_configured_model_name(config)
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
467 if no_documents:
468 logger.warning("No documents found to index.")
470 chunk_size = config.llm_embedding_chunk_size
471 embed_model = get_embedding_model(config)
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 )
508 _with_exclusive_access("compaction", store.compact)
509 return msg
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))
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)
543def llm_index_migrate() -> None:
544 """Apply any pending LLM index schema migrations, with no reindex.
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 )
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))
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))
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
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)
633def truncate_embedding_query(content: str, *, chunk_size: int) -> str:
634 from llama_index.core.text_splitter import TokenTextSplitter
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 ""
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}
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 []
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 []
672 config = AIConfig()
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
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)
684 filters = (
685 MetadataFilters(filters=filter_parts, condition=FilterCondition.AND)
686 if filter_parts
687 else None
688 )
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)
709 if allowed_document_ids is None:
710 return results
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