diff --git a/src/documents/apps.py b/src/documents/apps.py index c14f56ee4..c4976b671 100644 --- a/src/documents/apps.py +++ b/src/documents/apps.py @@ -31,6 +31,7 @@ class DocumentsConfig(AppConfig): document_consumption_finished.connect(add_or_update_document_in_llm_index) document_updated.connect(run_workflows_updated) document_updated.connect(send_websocket_document_updated) + document_updated.connect(add_or_update_document_in_llm_index) import documents.schema # noqa: F401 diff --git a/src/paperless_ai/indexing.py b/src/paperless_ai/indexing.py index 5e2c5e369..33061d383 100644 --- a/src/paperless_ai/indexing.py +++ b/src/paperless_ai/indexing.py @@ -1,5 +1,6 @@ import logging import shutil +from collections import defaultdict from collections.abc import Iterable from datetime import timedelta from pathlib import Path @@ -7,6 +8,7 @@ from typing import TYPE_CHECKING from django.conf import settings from django.utils import timezone +from filelock import FileLock from documents.models import Document from documents.models import PaperlessTask @@ -28,7 +30,17 @@ RAG_NUM_OUTPUT = 512 RAG_CHUNK_OVERLAP = 200 +def _index_lock_path() -> Path: + """Return the path used as the file lock for FAISS index mutations.""" + return settings.LLM_INDEX_DIR / "index.lock" + + def queue_llm_index_update_if_needed(*, rebuild: bool, reason: str) -> bool: + # NOTE: The check-then-enqueue sequence below is non-atomic (TOCTOU): two + # concurrent workers can both observe no running task and both enqueue a + # full rebuild. This is wasteful but not data-corrupting — update_llm_index + # is itself protected by _index_lock_path(), so only one rebuild runs at a + # time and the second one is serialised after the first completes. from documents.tasks import llmindex_index has_running = PaperlessTask.objects.filter( @@ -237,73 +249,73 @@ def update_llm_index( documents = Document.objects.all() if not documents.exists(): - msg = "No documents found to index." - logger.warning(msg) - return msg + logger.warning("No documents found to index.") + if not rebuild and not vector_store_file_exists(): + return "No documents found to index." config = AIConfig() chunk_size = config.llm_embedding_chunk_size - if rebuild or not vector_store_file_exists(): - # remove meta.json to force re-detection of embedding dim - (settings.LLM_INDEX_DIR / "meta.json").unlink(missing_ok=True) - # Rebuild index from scratch - logger.info("Rebuilding LLM index.") - import llama_index.core.settings as llama_settings + with FileLock(_index_lock_path()): + if rebuild or not vector_store_file_exists(): + # remove meta.json to force re-detection of embedding dim + (settings.LLM_INDEX_DIR / "meta.json").unlink(missing_ok=True) + # Rebuild index from scratch + logger.info("Rebuilding LLM index.") + import llama_index.core.settings as llama_settings - embed_model = get_embedding_model() - llama_settings.Settings.embed_model = embed_model - storage_context = get_or_create_storage_context(rebuild=True) - for document in iter_wrapper(documents): - document_nodes = build_document_node(document, chunk_size=chunk_size) - nodes.extend(document_nodes) + embed_model = get_embedding_model() + llama_settings.Settings.embed_model = embed_model + storage_context = get_or_create_storage_context(rebuild=True) + for document in iter_wrapper(documents): + document_nodes = build_document_node(document, chunk_size=chunk_size) + nodes.extend(document_nodes) - index = VectorStoreIndex( - nodes=nodes, - storage_context=storage_context, - embed_model=embed_model, - show_progress=False, - ) - msg = "LLM index rebuilt successfully." - else: - # Update existing index - index = load_or_build_index() - all_node_ids = list(index.docstore.docs.keys()) - existing_nodes = { - node.metadata.get("document_id"): node - for node in index.docstore.get_nodes(all_node_ids) - } - - for document in iter_wrapper(documents): - doc_id = str(document.id) - document_modified = document.modified.isoformat() - - if doc_id in existing_nodes: - node = existing_nodes[doc_id] - node_modified = node.metadata.get("modified") - - if node_modified == document_modified: - continue - - # Again, delete from docstore, FAISS IndexFlatL2 are append-only - index.docstore.delete_document(node.node_id) - nodes.extend(build_document_node(document, chunk_size=chunk_size)) - else: - # New document, add it - nodes.extend(build_document_node(document, chunk_size=chunk_size)) - - if nodes: - msg = "LLM index updated successfully." - logger.info( - "Updating %d nodes in LLM index.", - len(nodes), + index = VectorStoreIndex( + nodes=nodes, + storage_context=storage_context, + embed_model=embed_model, + show_progress=False, ) - index.insert_nodes(nodes) + msg = "LLM index rebuilt successfully." else: - msg = "No changes detected in LLM index." - logger.info(msg) + # Update existing index + index = load_or_build_index() + existing_nodes: defaultdict[str, list] = defaultdict(list) + for node in index.docstore.docs.values(): + doc_id = node.metadata.get("document_id") + if doc_id is not None: + existing_nodes[doc_id].append(node) - index.storage_context.persist(persist_dir=settings.LLM_INDEX_DIR) + for document in iter_wrapper(documents): + doc_id = str(document.id) + document_modified = document.modified.isoformat() + + if doc_id in existing_nodes: + doc_nodes = existing_nodes[doc_id] + node_modified = doc_nodes[0].metadata.get("modified") + + if node_modified == document_modified: + continue + + # Delete from docstore, FAISS IndexFlatL2 are append-only + for node in doc_nodes: + index.docstore.delete_document(node.node_id) + + nodes.extend(build_document_node(document, chunk_size=chunk_size)) + + if nodes: + msg = "LLM index updated successfully." + logger.info( + "Updating %d nodes in LLM index.", + len(nodes), + ) + index.insert_nodes(nodes) + else: + msg = "No changes detected in LLM index." + logger.info(msg) + + index.storage_context.persist(persist_dir=settings.LLM_INDEX_DIR) return msg @@ -313,25 +325,33 @@ def llm_index_add_or_update_document(document: Document): If the document already exists, it will be replaced. """ new_nodes = build_document_node(document, chunk_size=get_rag_chunk_size()) + if not new_nodes: + logger.warning( + "No indexable content for document %s; skipping LLM index update.", + document.pk, + ) + return - index = load_or_build_index(nodes=new_nodes) + with FileLock(_index_lock_path()): + index = load_or_build_index(nodes=new_nodes) - remove_document_docstore_nodes(document, index) + remove_document_docstore_nodes(document, index) - index.insert_nodes(new_nodes) + index.insert_nodes(new_nodes) - index.storage_context.persist(persist_dir=settings.LLM_INDEX_DIR) + index.storage_context.persist(persist_dir=settings.LLM_INDEX_DIR) def llm_index_remove_document(document: Document): """ Removes a document from the LLM index. """ - index = load_or_build_index() + with FileLock(_index_lock_path()): + index = load_or_build_index() - remove_document_docstore_nodes(document, index) + remove_document_docstore_nodes(document, index) - index.storage_context.persist(persist_dir=settings.LLM_INDEX_DIR) + index.storage_context.persist(persist_dir=settings.LLM_INDEX_DIR) def truncate_content( diff --git a/src/paperless_ai/matching.py b/src/paperless_ai/matching.py index f1dfc62db..c47c95001 100644 --- a/src/paperless_ai/matching.py +++ b/src/paperless_ai/matching.py @@ -98,5 +98,5 @@ def extract_unmatched_names( matched_objects: list, attr="name", ) -> list[str]: - matched_names = {getattr(obj, attr).lower() for obj in matched_objects} - return [name for name in names if name.lower() not in matched_names] + matched_names = {_normalize(getattr(obj, attr)) for obj in matched_objects} + return [name for name in names if _normalize(name) not in matched_names] diff --git a/src/paperless_ai/tests/test_ai_indexing.py b/src/paperless_ai/tests/test_ai_indexing.py index f0f66fb72..ee7fb4d67 100644 --- a/src/paperless_ai/tests/test_ai_indexing.py +++ b/src/paperless_ai/tests/test_ai_indexing.py @@ -1,15 +1,20 @@ import json +from pathlib import Path from unittest.mock import MagicMock from unittest.mock import patch import pytest +import pytest_mock from django.contrib.auth.models import User from django.test import override_settings from django.utils import timezone +from faker import Faker from llama_index.core.base.embeddings.base import BaseEmbedding from documents.models import Document from documents.models import PaperlessTask +from documents.signals import document_updated +from documents.tests.factories import DocumentFactory from documents.tests.factories import PaperlessTaskFactory from paperless.models import ApplicationConfiguration from paperless_ai import indexing @@ -505,6 +510,61 @@ def test_query_similar_documents_normalizes_and_post_filters_allowed_ids( assert private_document not in result +class TestUpdateLlmIndexStaleNodes: + """Tests that update_llm_index removes ALL nodes for a multi-chunk document.""" + + @pytest.mark.django_db + def test_incremental_update_removes_all_old_nodes_for_multi_chunk_document( + self, + temp_llm_index_dir, + mock_embed_model: MagicMock, + ) -> None: + """Ghost nodes from all chunks of a modified document must be removed. + + When a document is split into multiple chunks (chunk_size=1024), the + incremental update path must delete every old node, not just the last + one captured by a dict comprehension keyed on document_id. + """ + # Content long enough to produce at least two chunks at chunk_size=1024. + # Generate many paragraphs so the token count comfortably exceeds 1024. + fake = Faker() + long_content = "\n\n".join(fake.paragraph(nb_sentences=20) for _ in range(20)) + doc = DocumentFactory(content=long_content) + + # Build the initial index (rebuild=True) so it has multiple nodes + indexing.update_llm_index(rebuild=True) + + # Verify the initial index has more than one node for this document + initial_index = indexing.load_or_build_index() + initial_node_ids = [ + nid + for nid, node in initial_index.docstore.docs.items() + if node.metadata.get("document_id") == str(doc.id) + ] + assert len(initial_node_ids) > 1, ( + f"Expected multiple chunks but got {len(initial_node_ids)}; " + "increase long_content length" + ) + + # Simulate a modification so the incremental path treats it as changed. + # Use queryset.update() to bypass auto_now and actually change the DB value. + new_modified = timezone.now() + Document.objects.filter(pk=doc.pk).update(modified=new_modified) + + # Run incremental update (rebuild=False) with the modified document + indexing.update_llm_index(rebuild=False) + + # Reload the persisted index and check that no OLD node ids remain + updated_index = indexing.load_or_build_index() + remaining_old_node_ids = [ + nid for nid in initial_node_ids if nid in updated_index.docstore.docs + ] + assert remaining_old_node_ids == [], ( + f"Ghost nodes still present after incremental update: " + f"{remaining_old_node_ids}" + ) + + @pytest.mark.django_db def test_query_similar_documents_empty_allow_list_fails_closed( real_document, @@ -526,3 +586,220 @@ def test_query_similar_documents_empty_allow_list_fails_closed( mock_vector_store_exists.assert_not_called() mock_load_or_build_index.assert_not_called() mock_retriever_cls.assert_not_called() + + +class TestUpdateLlmIndexEmptyDocumentSet: + """update_llm_index must persist an empty index when all documents are deleted. + + Without this, the stale on-disk FAISS vectors are never cleared and + subsequent similarity searches return phantom hits for document IDs that + no longer exist in the DB. + """ + + @pytest.mark.django_db + def test_rebuild_clears_stale_index_when_no_documents_exist( + self, + temp_llm_index_dir: Path, + mock_embed_model: MagicMock, + ) -> None: + """After deleting all documents, rebuild=True must persist an empty index. + + Steps: + 1. Build an index with one document so the on-disk state is non-empty. + 2. Delete all documents from the DB. + 3. Call update_llm_index(rebuild=True). + 4. Reload the index from disk. + 5. Assert the reloaded index has zero nodes (no phantom vectors). + """ + # Step 1: create a document and build a non-empty index + Document.objects.create( + title="Soon-to-be-deleted document", + content="Some content that will become a phantom vector.", + added=timezone.now(), + ) + indexing.update_llm_index(rebuild=True) + + initial_index = indexing.load_or_build_index() + assert len(initial_index.docstore.docs) > 0, ( + "Precondition failed: expected at least one node before deletion" + ) + + # Step 2: delete all documents + Document.objects.all().delete() + assert not Document.objects.exists() + + # Step 3: rebuild with no documents + indexing.update_llm_index(rebuild=True) + + # Step 4: reload the persisted index from disk + reloaded_index = indexing.load_or_build_index() + + # Step 5: phantom vectors must be gone + assert len(reloaded_index.docstore.docs) == 0, ( + f"Expected 0 nodes after clearing all documents, " + f"but found {len(reloaded_index.docstore.docs)}: " + f"{list(reloaded_index.docstore.docs.keys())}" + ) + + +class TestDocumentUpdatedSignalTriggersLlmReindex: + """document_updated must enqueue an LLM index update, just like document_consumption_finished.""" + + @pytest.mark.django_db + @override_settings(AI_ENABLED=True, LLM_EMBEDDING_BACKEND="huggingface") + def test_document_updated_enqueues_llm_reindex( + self, + mocker: pytest_mock.MockerFixture, + ) -> None: + """Firing document_updated should call update_document_in_llm_index.apply_async.""" + mock_task = mocker.patch("documents.tasks.update_document_in_llm_index") + + doc = DocumentFactory() + document_updated.send(sender=object, document=doc) + + mock_task.apply_async.assert_called_once_with(kwargs={"document": doc}) + + +@pytest.mark.django_db +class TestLlmIndexAddOrUpdateDocumentEmptyContent: + """llm_index_add_or_update_document must handle empty node lists gracefully.""" + + def test_returns_without_error_when_build_document_node_returns_empty( + self, + temp_llm_index_dir: Path, + mocker: pytest_mock.MockerFixture, + ) -> None: + """When build_document_node returns [], the function must return without error + and must not call load_or_build_index at all.""" + mocker.patch( + "paperless_ai.indexing.build_document_node", + return_value=[], + ) + mock_load = mocker.patch("paperless_ai.indexing.load_or_build_index") + + doc = MagicMock(spec=Document) + # Must not raise + indexing.llm_index_add_or_update_document(doc) + + mock_load.assert_not_called() + + +@pytest.mark.django_db +class TestLlmIndexLocking: + """The FAISS index mutation functions must acquire the index lock before touching the index. + + Without locking, two concurrent Celery workers can each load the same + on-disk index, make independent modifications, and the last writer silently + overwrites the first's changes. + """ + + def test_add_or_update_document_acquires_lock( + self, + temp_llm_index_dir: Path, + mocker: pytest_mock.MockerFixture, + ) -> None: + """llm_index_add_or_update_document must enter the file lock before touching the index.""" + call_order: list[str] = [] + + mock_lock_instance = MagicMock() + mock_lock_instance.__enter__ = MagicMock( + side_effect=lambda *_: call_order.append("lock_acquired"), + ) + mock_lock_instance.__exit__ = MagicMock(return_value=False) + + mock_file_lock_cls = mocker.patch( + "paperless_ai.indexing.FileLock", + return_value=mock_lock_instance, + ) + + mock_load = mocker.patch( + "paperless_ai.indexing.load_or_build_index", + side_effect=lambda *_a, **_kw: ( + call_order.append("index_loaded") or MagicMock() + ), + ) + mocker.patch( + "paperless_ai.indexing.build_document_node", + return_value=[MagicMock()], + ) + mocker.patch("paperless_ai.indexing.remove_document_docstore_nodes") + + doc = MagicMock(spec=Document) + indexing.llm_index_add_or_update_document(doc) + + mock_file_lock_cls.assert_called_once() + mock_lock_instance.__enter__.assert_called_once() + mock_load.assert_called_once() + assert call_order.index("lock_acquired") < call_order.index("index_loaded"), ( + "Lock must be acquired before the index is loaded" + ) + + def test_remove_document_acquires_lock( + self, + temp_llm_index_dir: Path, + mocker: pytest_mock.MockerFixture, + ) -> None: + """llm_index_remove_document must enter the file lock before loading the index.""" + call_order: list[str] = [] + + mock_lock_instance = MagicMock() + mock_lock_instance.__enter__ = MagicMock( + side_effect=lambda *_: call_order.append("lock_acquired"), + ) + mock_lock_instance.__exit__ = MagicMock(return_value=False) + + mock_file_lock_cls = mocker.patch( + "paperless_ai.indexing.FileLock", + return_value=mock_lock_instance, + ) + + mock_load = mocker.patch( + "paperless_ai.indexing.load_or_build_index", + side_effect=lambda *_a, **_kw: ( + call_order.append("index_loaded") or MagicMock() + ), + ) + mocker.patch("paperless_ai.indexing.remove_document_docstore_nodes") + + doc = MagicMock(spec=Document) + indexing.llm_index_remove_document(doc) + + mock_file_lock_cls.assert_called_once() + mock_lock_instance.__enter__.assert_called_once() + mock_load.assert_called_once() + assert call_order.index("lock_acquired") < call_order.index("index_loaded"), ( + "Lock must be acquired before the index is loaded" + ) + + def test_update_llm_index_rebuild_acquires_lock( + self, + temp_llm_index_dir: Path, + mock_embed_model: MagicMock, + mocker: pytest_mock.MockerFixture, + ) -> None: + """update_llm_index must enter the file lock during the rebuild/persist cycle.""" + mock_lock_instance = MagicMock() + mock_lock_instance.__enter__ = MagicMock(return_value=None) + mock_lock_instance.__exit__ = MagicMock(return_value=False) + + mock_file_lock_cls = mocker.patch( + "paperless_ai.indexing.FileLock", + return_value=mock_lock_instance, + ) + + # exists=True so the code reaches the lock; iterate over an empty + # queryset so VectorStoreIndex is called with no nodes (still exercises + # the lock path without needing heavy FAISS fixture data) + mock_qs = MagicMock() + mock_qs.exists.return_value = True + mock_qs.__iter__ = MagicMock(return_value=iter([])) + mocker.patch("paperless_ai.indexing.Document.objects.all", return_value=mock_qs) + mocker.patch( + "paperless_ai.indexing.get_or_create_storage_context", + return_value=MagicMock(), + ) + + indexing.update_llm_index(rebuild=True) + + mock_file_lock_cls.assert_called_once() + mock_lock_instance.__enter__.assert_called_once() diff --git a/src/paperless_ai/tests/test_matching.py b/src/paperless_ai/tests/test_matching.py index 83cfd8a41..5cf23f2b8 100644 --- a/src/paperless_ai/tests/test_matching.py +++ b/src/paperless_ai/tests/test_matching.py @@ -1,5 +1,6 @@ from unittest.mock import patch +import pytest from django.test import TestCase from documents.models import Correspondent @@ -84,3 +85,17 @@ class TestAIMatching(TestCase): self.assertEqual(len(result), 2) self.assertEqual(result[0].name, "Test Tag 1") self.assertEqual(result[1].name, "Test Tag 2") + + +@pytest.mark.django_db +class TestExtractUnmatchedNamesNormalization: + def test_punctuated_name_already_matched_is_not_returned_as_unmatched( + self, + ) -> None: + correspondent = Correspondent.objects.create(name="J Smith") + llm_names = ["J. Smith"] + matched_objects: list[Correspondent] = [correspondent] + + unmatched = extract_unmatched_names(llm_names, matched_objects) + + assert "J. Smith" not in unmatched