Compare commits

..
22 changed files with 519 additions and 328 deletions
@@ -0,0 +1,12 @@
#!/command/with-contenv /usr/bin/bash
# shellcheck shell=bash
declare -r log_prefix="[init-llmindex-migrate]"
echo "${log_prefix} Checking for pending LLM index migrations..."
cd "${PAPERLESS_SRC_DIR}"
if [[ -n "${USER_IS_NON_ROOT}" ]]; then
python3 manage.py document_llmindex migrate
else
s6-setuidgid paperless python3 manage.py document_llmindex migrate
fi
@@ -0,0 +1 @@
oneshot
@@ -0,0 +1 @@
/etc/s6-overlay/s6-rc.d/init-llmindex-migrate/run
+11 -1
View File
@@ -212,6 +212,16 @@ following:
This is a no-op if the index is already up to date, so it is safe to
run on every upgrade.
5. Migrate the LLM index if needed.
```shell-session
cd src
python3 manage.py document_llmindex migrate
```
This is a no-op if the index schema is already current, so it is safe
to run on every upgrade.
### Database Upgrades
Paperless-ngx is compatible with Django-supported versions of PostgreSQL and MariaDB and it is generally
@@ -532,7 +542,7 @@ index is updated automatically on the schedule set by
can manage it manually:
```
document_llmindex {rebuild,update,compact}
document_llmindex {rebuild,update,compact,migrate}
```
Specify `rebuild` to build the index from scratch from all documents in the database. Use
+25 -15
View File
@@ -57,7 +57,9 @@ from paperless.models import ArchiveFileGenerationChoices
from paperless.parsers import ParserContext
from paperless.parsers import ParserProtocol
from paperless.parsers.registry import get_parser_registry
from paperless.parsers.utils import pdf_born_digital_text
from paperless.parsers.utils import PDF_TEXT_MIN_LENGTH
from paperless.parsers.utils import extract_pdf_text
from paperless.parsers.utils import is_tagged_pdf
LOGGING_NAME: Final[str] = "paperless.consumer"
@@ -136,45 +138,53 @@ def should_produce_archive(
# Must produce a PDF so the frontend can display the original format at all.
if parser.requires_pdf_rendition:
_log.debug("Archive: yes - parser requires PDF rendition for frontend display")
_log.debug("Archive: yes parser requires PDF rendition for frontend display")
return True
# Parser cannot produce an archive (e.g. TextDocumentParser).
if not parser.can_produce_archive:
_log.debug("Archive: no - parser cannot produce archives")
_log.debug("Archive: no parser cannot produce archives")
return False
generation = OcrConfig().archive_file_generation
if generation == ArchiveFileGenerationChoices.ALWAYS:
_log.debug("Archive: yes - ARCHIVE_FILE_GENERATION=always")
_log.debug("Archive: yes ARCHIVE_FILE_GENERATION=always")
return True
if generation == ArchiveFileGenerationChoices.NEVER:
_log.debug("Archive: no - ARCHIVE_FILE_GENERATION=never")
_log.debug("Archive: no ARCHIVE_FILE_GENERATION=never")
return False
# auto: produce archives for scanned/image documents; skip for born-digital PDFs.
if mime_type.startswith("image/"):
_log.debug("Archive: yes - image document, ARCHIVE_FILE_GENERATION=auto")
_log.debug("Archive: yes image document, ARCHIVE_FILE_GENERATION=auto")
return True
if mime_type == "application/pdf":
text, born_digital = pdf_born_digital_text(document_path, log=_log)
text_length = len(text) if text else 0
if born_digital:
text = extract_pdf_text(document_path)
has_text = text is not None and len(text) > 0
if has_text and is_tagged_pdf(document_path):
_log.debug(
"Archive: no - born-digital PDF (text_length=%d),"
"Archive: no born-digital PDF (structure tags detected),"
" ARCHIVE_FILE_GENERATION=auto",
text_length,
)
return False
if text is None or len(text) <= PDF_TEXT_MIN_LENGTH:
_log.debug(
"Archive: yes — scanned PDF (text_length=%d%d),"
" ARCHIVE_FILE_GENERATION=auto",
len(text) if text else 0,
PDF_TEXT_MIN_LENGTH,
)
return True
_log.debug(
"Archive: yes - scanned/textless PDF (text_length=%d),"
"Archive: no — born-digital PDF (text_length=%d > %d),"
" ARCHIVE_FILE_GENERATION=auto",
text_length,
len(text),
PDF_TEXT_MIN_LENGTH,
)
return True
return False
_log.debug(
"Archive: no - MIME type %r not eligible for auto archive generation",
"Archive: no MIME type %r not eligible for auto archive generation",
mime_type,
)
return False
@@ -3,6 +3,7 @@ from typing import Any
from documents.management.commands.base import PaperlessCommand
from documents.tasks import llmindex_index
from paperless_ai.indexing import llm_index_compact
from paperless_ai.indexing import llm_index_migrate
class Command(PaperlessCommand):
@@ -13,12 +14,18 @@ class Command(PaperlessCommand):
def add_arguments(self, parser: Any) -> None:
super().add_arguments(parser)
parser.add_argument("command", choices=["rebuild", "update", "compact"])
parser.add_argument(
"command",
choices=["rebuild", "update", "compact", "migrate"],
)
def handle(self, *args: Any, **options: Any) -> None:
if options["command"] == "compact":
llm_index_compact()
return
if options["command"] == "migrate":
llm_index_migrate()
return
llmindex_index(
rebuild=options["command"] == "rebuild",
iter_wrapper=lambda docs: self.track(
@@ -9,6 +9,7 @@ if TYPE_CHECKING:
_COMPACT = "documents.management.commands.document_llmindex.llm_index_compact"
_INDEX = "documents.management.commands.document_llmindex.llmindex_index"
_MIGRATE = "documents.management.commands.document_llmindex.llm_index_migrate"
class TestDocumentLlmindexCommand:
@@ -17,6 +18,11 @@ class TestDocumentLlmindexCommand:
call_command("document_llmindex", "compact")
mock_compact.assert_called_once_with()
def test_migrate_calls_llm_index_migrate(self, mocker: MockerFixture) -> None:
mock_migrate = mocker.patch(_MIGRATE)
call_command("document_llmindex", "migrate")
mock_migrate.assert_called_once_with()
def test_rebuild_calls_llmindex_index_with_rebuild_true(
self,
mocker: MockerFixture,
+2 -2
View File
@@ -1329,7 +1329,7 @@ class PreConsumeTestCase(DirectoriesMixin, GetConsumerMixin, TestCase):
with self.get_consumer(self.test_file) as c:
c.run()
# Verify no pre-consume script subprocess was invoked
# (run_subprocess may still be called by pdf_born_digital_text via pdftotext)
# (run_subprocess may still be called by _extract_text_for_archive_check)
script_calls = [
call
for call in m.call_args_list
@@ -1354,7 +1354,7 @@ class PreConsumeTestCase(DirectoriesMixin, GetConsumerMixin, TestCase):
self.assertTrue(m.called)
# Find the call that invoked the pre-consume script
# (run_subprocess may also be called by pdf_born_digital_text via pdftotext)
# (run_subprocess may also be called by _extract_text_for_archive_check)
script_call = next(
call
for call in m.call_args_list
+42 -14
View File
@@ -134,32 +134,60 @@ class TestShouldProduceArchive:
assert should_produce_archive(parser, mime, Path("/tmp/doc")) is expected
@pytest.mark.parametrize(
("born_digital", "expected"),
("extracted_text", "expected"),
[
pytest.param(True, False, id="born-digital-skips-archive"),
pytest.param(False, True, id="not-born-digital-produces-archive"),
pytest.param(
"This is a born-digital PDF with lots of text content. " * 10,
False,
id="born-digital-long-text-skips-archive",
),
pytest.param(None, True, id="no-text-scanned-produces-archive"),
pytest.param("tiny", True, id="short-text-treated-as-scanned"),
],
)
def test_auto_pdf_archive_decision(
self,
mocker: MockerFixture,
settings,
born_digital: bool, # noqa: FBT001
extracted_text: str | None,
expected: bool, # noqa: FBT001
) -> None:
"""Archive decision tracks pdf_born_digital_text()'s verdict exactly.
should_produce_archive() defers entirely to pdf_born_digital_text()
for the has-real-text decision, so both callers of that predicate
(this function and RasterisedDocumentParser.parse()) always agree.
"""
settings.ARCHIVE_FILE_GENERATION = "auto"
mocker.patch(
"documents.consumer.pdf_born_digital_text",
return_value=("some text", born_digital),
)
mocker.patch("documents.consumer.is_tagged_pdf", return_value=False)
mocker.patch("documents.consumer.extract_pdf_text", return_value=extracted_text)
parser = _parser_instance(can_produce=True, requires_rendition=False)
assert (
should_produce_archive(parser, "application/pdf", Path("/tmp/doc.pdf"))
is expected
)
def test_tagged_pdf_skips_archive_in_auto_mode(
self,
mocker: MockerFixture,
settings,
) -> None:
"""Tagged PDFs (e.g. Word exports) with real text are treated as born-digital, even below PDF_TEXT_MIN_LENGTH."""
settings.ARCHIVE_FILE_GENERATION = "auto"
mocker.patch("documents.consumer.is_tagged_pdf", return_value=True)
mocker.patch("documents.consumer.extract_pdf_text", return_value="tiny")
parser = _parser_instance(can_produce=True, requires_rendition=False)
assert (
should_produce_archive(parser, "application/pdf", Path("/tmp/doc.pdf"))
is False
)
def test_tagged_pdf_without_text_produces_archive(
self,
mocker: MockerFixture,
settings,
) -> None:
"""A tagged PDF with no actual extractable text (e.g. some scanner firmware) is not
trusted as born-digital — the tag alone must not bypass OCR."""
settings.ARCHIVE_FILE_GENERATION = "auto"
mocker.patch("documents.consumer.is_tagged_pdf", return_value=True)
mocker.patch("documents.consumer.extract_pdf_text", return_value=None)
parser = _parser_instance(can_produce=True, requires_rendition=False)
assert (
should_produce_archive(parser, "application/pdf", Path("/tmp/doc.pdf"))
is True
)
+21 -6
View File
@@ -3,6 +3,7 @@ from __future__ import annotations
import importlib.resources
import logging
import os
import re
import shutil
import tempfile
from pathlib import Path
@@ -24,9 +25,9 @@ from paperless.config import OcrConfig
from paperless.models import CleanChoices
from paperless.models import ModeChoices
from paperless.models import OutputTypeChoices
from paperless.parsers.utils import PDF_TEXT_MIN_LENGTH
from paperless.parsers.utils import extract_pdf_text
from paperless.parsers.utils import is_born_digital_text
from paperless.parsers.utils import post_process_text
from paperless.parsers.utils import is_tagged_pdf
from paperless.parsers.utils import read_file_handle_unicode_errors
from paperless.version import __full_version_str__
@@ -509,10 +510,10 @@ class RasterisedDocumentParser:
if mime_type == "application/pdf":
text_original = self.extract_text(None, document_path)
original_has_text = is_born_digital_text(
text_original,
document_path,
log=self.log,
has_text = text_original is not None and len(text_original) > 0
original_has_text = has_text and (
is_tagged_pdf(document_path, log=self.log)
or len(text_original) > PDF_TEXT_MIN_LENGTH
)
else:
text_original = None
@@ -657,3 +658,17 @@ class RasterisedDocumentParser:
f"No text was found in {document_path}, the content will be empty.",
)
self.text = ""
def post_process_text(text: str | None) -> str | None:
if not text:
return None
collapsed_spaces = re.sub(r"([^\S\r\n]+)", " ", text)
no_leading_whitespace = re.sub(r"([\n\r]+)([^\S\n\r]+)", "\\1", collapsed_spaces)
no_trailing_whitespace = re.sub(r"([^\S\n\r]+)$", "", no_leading_whitespace)
# TODO: this needs a rework
# replace \0 prevents issues with saving to postgres.
# text may contain \0 when this character is present in PDF files.
return no_trailing_whitespace.strip().replace("\0", " ")
-82
View File
@@ -111,88 +111,6 @@ def extract_pdf_text(
return None
def post_process_text(text: str | None) -> str | None:
"""Normalize extracted PDF/OCR text: collapse whitespace, strip padding.
Returns ``None`` for ``None`` or whitespace-only input, so callers can
treat "no text" and "only layout padding" the same way.
"""
if not text:
return None
collapsed_spaces = re.sub(r"([^\S\r\n]+)", " ", text)
no_leading_whitespace = re.sub(r"([\n\r]+)([^\S\n\r]+)", "\\1", collapsed_spaces)
no_trailing_whitespace = re.sub(r"([^\S\n\r]+)$", "", no_leading_whitespace)
# replace \0 prevents issues with saving to postgres.
# text may contain \0 when this character is present in PDF files.
result = no_trailing_whitespace.strip().replace("\0", " ")
return result or None
def is_born_digital_text(
text: str | None,
path: Path,
log: logging.Logger | None = None,
) -> bool:
"""Decide whether already-extracted, normalized PDF text counts as born-digital.
This is the single source of truth for "does this PDF already have real
text", used both to decide whether to produce an archive file and to
decide whether OCR can be skipped. Both decisions must agree, or a
tagged-but-textless PDF can end up with no archive AND a forced OCR pass
(see GH #13387): raw ``pdftotext -layout`` output can be non-empty
(whitespace/form-feed padding) even when there is no real content, so
*text* must already be normalized via :func:`post_process_text`, not the
raw extraction.
Parameters
----------
text:
The normalized extracted text (or ``None``) to evaluate.
path:
Absolute path to the PDF file, used for the tagged-PDF check.
log:
Logger for warnings. Falls back to the module-level logger when omitted.
Returns
-------
bool
Whether the PDF counts as born-digital (has real text, and is either
tagged or exceeds ``PDF_TEXT_MIN_LENGTH``).
"""
if not text:
return False
return is_tagged_pdf(path, log=log) or len(text) > PDF_TEXT_MIN_LENGTH
def pdf_born_digital_text(
path: Path,
log: logging.Logger | None = None,
) -> tuple[str | None, bool]:
"""Extract a PDF's text and decide whether it should be treated as born-digital.
Convenience wrapper around :func:`is_born_digital_text` for callers that
don't already have the PDF's text extracted (e.g. the archive-generation
decision, which runs before any parser has touched the file).
Parameters
----------
path:
Absolute path to the PDF file.
log:
Logger for warnings. Falls back to the module-level logger when omitted.
Returns
-------
tuple[str | None, bool]
The normalized extracted text (or ``None``), and whether the PDF
counts as born-digital.
"""
text = post_process_text(extract_pdf_text(path, log=log))
return text, is_born_digital_text(text, path, log=log)
def read_file_handle_unicode_errors(
filepath: Path,
log: logging.Logger | None = None,
-17
View File
@@ -36,23 +36,6 @@ def samples_dir() -> Path:
return (Path(__file__).parent / "samples").resolve()
@pytest.fixture(scope="session")
def tagged_no_text_pdf_file(samples_dir: Path) -> Path:
"""Path to a tagged PDF whose only "text" is pdftotext layout padding.
Reproduces GH #13387: ``/MarkInfo /Marked true`` is set, but the only
extractable content is a form-feed byte, not real text. Lives here
rather than in parsers/conftest.py so both parser tests and
paperless/tests/test_parser_utils.py can use it.
Returns
-------
Path
Absolute path to ``tesseract/tagged-but-no-text.pdf``.
"""
return samples_dir / "tesseract" / "tagged-but-no-text.pdf"
@pytest.fixture(autouse=True)
def clean_registry() -> Generator[None, None, None]:
"""Reset the parser registry before and after every test.
@@ -21,7 +21,7 @@ from documents.parsers import run_convert
from paperless.models import ModeChoices
from paperless.parsers import ParserProtocol
from paperless.parsers.tesseract import RasterisedDocumentParser
from paperless.parsers.utils import is_tagged_pdf
from paperless.parsers.tesseract import post_process_text
if TYPE_CHECKING:
from pathlib import Path
@@ -151,6 +151,36 @@ class TestRasterisedDocumentParserLifecycle:
assert tempdir is not None and not tempdir.exists()
# ---------------------------------------------------------------------------
# post_process_text
# ---------------------------------------------------------------------------
class TestPostProcessText:
@pytest.mark.parametrize(
("source", "expected"),
[
pytest.param(
"simple string",
"simple string",
id="collapse-spaces",
),
pytest.param(
"simple newline\n testing string",
"simple newline\ntesting string",
id="preserve-newline",
),
pytest.param(
"utf-8 строка с пробелами в конце ", # noqa: RUF001
"utf-8 строка с пробелами в конце", # noqa: RUF001
id="utf8-trailing-spaces",
),
],
)
def test_post_process_text(self, source: str, expected: str) -> None:
assert post_process_text(source) == expected
# ---------------------------------------------------------------------------
# Page count
# ---------------------------------------------------------------------------
@@ -880,25 +910,25 @@ class TestSkipArchive:
self,
mocker: MockerFixture,
tesseract_parser: RasterisedDocumentParser,
tagged_no_text_pdf_file: Path,
tesseract_samples_dir: Path,
) -> None:
"""
GIVEN:
- A real PDF that reports itself as tagged (/MarkInfo /Marked
true) but whose only pdftotext output is layout padding (a
lone form-feed byte), not real text (see GitHub issue #13387,
originally reported against #13349's tagged-PDF handling)
- A PDF that reports itself as tagged (/MarkInfo /Marked true) but
has no actual extractable text (some scanner firmware produces
this — see GitHub issue #13349)
- Mode: auto, produce_archive=False
WHEN:
- Document is parsed
THEN:
- The tag alone is not trusted as "has text"; OCRmyPDF still runs
"""
assert is_tagged_pdf(tagged_no_text_pdf_file) is True
tesseract_parser.settings.mode = ModeChoices.AUTO
mocker.patch("paperless.parsers.tesseract.is_tagged_pdf", return_value=True)
mocker.patch.object(tesseract_parser, "extract_text", return_value=None)
mock_ocr = mocker.patch("ocrmypdf.ocr")
tesseract_parser.parse(
tagged_no_text_pdf_file,
tesseract_samples_dir / "multi-page-images.pdf",
"application/pdf",
produce_archive=False,
)
-110
View File
@@ -4,18 +4,10 @@ from __future__ import annotations
import codecs
from pathlib import Path
from typing import TYPE_CHECKING
import pytest
from paperless.parsers.utils import is_tagged_pdf
from paperless.parsers.utils import pdf_born_digital_text
from paperless.parsers.utils import post_process_text
from paperless.parsers.utils import read_file_handle_unicode_errors
if TYPE_CHECKING:
from pytest_mock import MockerFixture
SAMPLES = Path(__file__).parent / "samples" / "tesseract"
@@ -68,105 +60,3 @@ class TestIsTaggedPdf:
bad = tmp_path / "bad.pdf"
bad.write_bytes(b"not a pdf")
assert is_tagged_pdf(bad) is False
class TestPostProcessText:
@pytest.mark.parametrize(
("source", "expected"),
[
pytest.param(
"simple string",
"simple string",
id="collapse-spaces",
),
pytest.param(
"simple newline\n testing string",
"simple newline\ntesting string",
id="preserve-newline",
),
pytest.param(
"utf-8 строка с пробелами в конце ", # noqa: RUF001
"utf-8 строка с пробелами в конце", # noqa: RUF001
id="utf8-trailing-spaces",
),
pytest.param(None, None, id="none-input"),
pytest.param("", None, id="empty-string"),
pytest.param(" \n\x0c \n ", None, id="whitespace-and-formfeed-only"),
],
)
def test_post_process_text(
self,
source: str | None,
expected: str | None,
) -> None:
assert post_process_text(source) == expected
class TestPdfBornDigitalText:
"""Regression coverage for GH #13387.
should_produce_archive() and RasterisedDocumentParser.parse() must agree
on whether a PDF has real text, so both go through this one function.
"""
@pytest.mark.parametrize(
("extracted", "tagged", "expected_text", "expected_born_digital"),
[
pytest.param("tiny", True, "tiny", True, id="tagged-with-real-text"),
pytest.param("tiny", False, "tiny", False, id="untagged-below-min-length"),
pytest.param(
"x" * 51,
False,
"x" * 51,
True,
id="untagged-above-min-length",
),
pytest.param(None, True, None, False, id="tagged-but-no-text"),
],
)
def test_born_digital_decision(
self,
mocker: MockerFixture,
tmp_path: Path,
extracted: str | None,
tagged: bool, # noqa: FBT001
expected_text: str | None,
expected_born_digital: bool, # noqa: FBT001
) -> None:
"""
GIVEN:
- A PDF whose pdftotext output and /MarkInfo tag status vary
WHEN:
- pdf_born_digital_text() is called
THEN:
- The normalized text and born-digital verdict match; the tag
alone never counts as "has text"
"""
mocker.patch(
"paperless.parsers.utils.extract_pdf_text",
return_value=extracted,
)
mocker.patch("paperless.parsers.utils.is_tagged_pdf", return_value=tagged)
text, born_digital = pdf_born_digital_text(tmp_path / "doc.pdf")
assert text == expected_text
assert born_digital is expected_born_digital
def test_tagged_but_textless_pdf_is_not_born_digital(
self,
tagged_no_text_pdf_file: Path,
) -> None:
"""
GIVEN:
- A real PDF that is tagged (/MarkInfo /Marked true) but whose
only "text" is layout padding (a stray form-feed byte)
WHEN:
- pdf_born_digital_text() is called with no mocking
THEN:
- The normalized text is None and the PDF is not treated as
born-digital. The raw, unnormalized pdftotext output is
non-empty for this file, which is exactly what caused the
archive decision to disagree with the OCR decision in #13387.
"""
text, born_digital = pdf_born_digital_text(tagged_no_text_pdf_file)
assert text is None
assert born_digital is False
+78 -25
View File
@@ -144,6 +144,24 @@ def _exclude_readers():
lock.close()
def _with_exclusive_access(operation: str, fn):
"""Run ``fn()`` with exclusive index access (see ``_exclude_readers()``),
for compaction/migration file swaps that must not run while readers are
active. Returns ``fn()``'s result, or None (after logging) if active
readers do not drain within ``LLM_INDEX_COMPACTION_LOCK_TIMEOUT`` --
callers skip the operation this run; it retries next time.
"""
try:
with _exclude_readers():
return fn()
except Timeout:
logger.info(
"Skipping LLM index %s: index readers are active; will retry next run.",
operation,
)
return None
@contextmanager
def write_store(embed_model_name: str | None = None):
"""Acquire the write lock and yield the vector store.
@@ -168,6 +186,21 @@ def write_store(embed_model_name: str | None = None):
yield store
def _check_and_run_migrations(store: "PaperlessSqliteVecVectorStore") -> bool:
"""Run any pending structural migrations, returning True if a pending
re-embed migration needs the caller to force a rebuild -- never
triggered automatically here. Safe to call before any write, including
delete()/upsert_document(): has_pending_migration() (see its docstring)
keeps this a no-op, with no exclusive access taken, once the store is
current.
"""
if not store.has_pending_migration():
return False
return bool(
_with_exclusive_access("migration check", store.check_and_run_migrations),
)
def _safe_related_name(document: Document, field: str) -> str | None:
"""
Returns the ``name`` of a related object (correspondent, document_type,
@@ -339,15 +372,7 @@ def update_llm_index(
happens, since a rebuild always covers the whole library regardless.
"""
with write_store() as store:
try:
with _exclude_readers():
needs_reembed = store.check_and_run_migrations()
except Timeout:
logger.info(
"Skipping LLM index migration check: index readers are active; "
"will retry next run.",
)
needs_reembed = False
needs_reembed = _check_and_run_migrations(store)
if needs_reembed:
logger.warning(
"LLM index migration requires re-embedding; forcing rebuild.",
@@ -412,14 +437,7 @@ def update_llm_index(
else "No changes detected in LLM index."
)
try:
with _exclude_readers():
store.compact()
except Timeout:
logger.info(
"Skipping LLM index compaction: index readers are active; "
"will retry next run.",
)
_with_exclusive_access("compaction", store.compact)
return msg
@@ -434,25 +452,60 @@ def llm_index_add_or_update_document(document: Document):
_embed_nodes(new_nodes, get_embedding_model(config))
with write_store(embed_model_name=get_configured_model_name(config)) as store:
needs_reembed = _check_and_run_migrations(store)
if needs_reembed:
logger.warning(
"Skipping incremental LLM index update for document %s: the "
"index requires re-embedding first. Run 'document_llmindex "
"rebuild' to resolve.",
document.id,
)
return
store.upsert_document(str(document.id), new_nodes)
def llm_index_migrate() -> None:
"""Apply any pending LLM index schema migrations, with no reindex.
Intended to run unconditionally on every startup (see the
init-llmindex-migrate container step and the bare-metal upgrade docs):
has_pending_migration() short-circuits to a metadata-only read once the
store is current, so a healthy install pays almost nothing here. Only
ever applies structural migrations -- a pending re-embed migration is
left for the explicit, deliberate rebuild path (``document_llmindex
update``/``rebuild``) to resolve, since re-embedding can be slow and,
for a metered embedding backend, cost money.
"""
if not AIConfig().llm_index_enabled:
return
with write_store() as store:
needs_reembed = _check_and_run_migrations(store)
if needs_reembed:
logger.warning(
"LLM index requires re-embedding, which this automatic migration "
"check will not do on its own -- it can be slow and, for a "
"metered embedding backend, cost money. Run "
"'document_llmindex rebuild' manually when ready.",
)
def llm_index_compact() -> None:
"""Compact the index immediately, rebuilding the table to reclaim space."""
with write_store() as store:
try:
with _exclude_readers():
store.compact(force=True)
except Timeout:
logger.info(
"Skipping LLM index compaction: index readers are active; "
"will retry next run.",
)
_with_exclusive_access("compaction", lambda: store.compact(force=True))
def llm_index_remove_document(document: Document):
"""Remove a document's chunks from the LLM index."""
with write_store() as store:
if _check_and_run_migrations(store):
logger.warning(
"Skipping removal of document %s from the LLM index: the "
"index requires re-embedding first. Run 'document_llmindex "
"rebuild' to resolve.",
document.id,
)
return
store.delete(str(document.id))
+60
View File
@@ -0,0 +1,60 @@
"""Schema migrations for the sqlite-vec vector store.
Each migration lives in its own module here, named ``mNNNN_description.py``
(e.g. ``m0001_v1_to_v2.py`` -- a leading digit isn't a valid Python
identifier, hence the ``m`` prefix, unlike Django's own numbered migrations,
which load via a dynamic ``importlib.import_module()`` call rather than a
static import statement), and registers itself into ``MIGRATIONS`` at import
time. ``vector_store.py`` imports those modules at the bottom of the file,
purely for that registration side effect, after ``PaperlessSqliteVecVectorStore``
is fully defined -- migrations need it to implement ``apply()`` (see
``Migration`` below).
To add a new migration: add a new ``mNNNN_description.py`` module here that
imports ``PaperlessSqliteVecVectorStore`` from ``paperless_ai.vector_store``,
defines its ``apply()``, and appends a ``Migration`` to ``MIGRATIONS``; then
import that module at the bottom of ``vector_store.py`` and bump
``SCHEMA_VERSION`` there. A migration must freeze its own historical DDL for
any side table its target version depends on (``DROP TABLE IF EXISTS`` +
its own literal ``CREATE TABLE``/``CREATE INDEX`` statements) rather than
delegating to any "current schema" helper -- see ``m0001_v1_to_v2.py`` for
why and the worked example.
"""
import sqlite3
from collections.abc import Callable
from dataclasses import dataclass
from dataclasses import field
from typing import Literal
@dataclass
class Migration:
"""A schema migration for the sqlite-vec vector store.
kind="structural": rows are copied into a new-schema file with no
re-embedding needed. Supply ``apply(src_conn, dst_conn, dim)``, which
must create every table its target schema needs in ``dst_conn`` and copy
``src_conn``'s rows and relevant ``index_meta`` keys into it.
``schema_version`` is written by the migration runner after ``apply``
returns, not by ``apply`` itself.
kind="re-embed": the new schema requires fresh embeddings.
``check_and_run_migrations()`` returns True when it encounters one of
these so the caller can force a full rebuild (which recreates the table
at the current SCHEMA_VERSION).
"""
from_version: int
to_version: int
kind: Literal["structural", "re-embed"]
description: str
apply: Callable[[sqlite3.Connection, sqlite3.Connection, int], None] | None = field(
default=None,
repr=False,
)
# Registry of all schema migrations in order, populated by each migration
# module's import-time registration (see the module docstring above).
MIGRATIONS: list[Migration] = []
+130
View File
@@ -1,3 +1,4 @@
import logging
from pathlib import Path
from unittest.mock import MagicMock
from unittest.mock import patch
@@ -737,6 +738,7 @@ class TestLlmIndexLocking:
mocker: pytest_mock.MockerFixture,
) -> None:
mock_store = MagicMock()
mock_store.has_pending_migration.return_value = False
mocker.patch(
"paperless_ai.indexing.write_store",
return_value=mocker.MagicMock(
@@ -757,12 +759,45 @@ class TestLlmIndexLocking:
mock_store.upsert_document.assert_called_once()
def test_add_or_update_document_skips_write_when_reembed_pending(
self,
temp_llm_index_dir: Path,
mock_embed_model: FakeEmbedding,
mocker: pytest_mock.MockerFixture,
) -> None:
"""A pending re-embed migration must block the incremental write,
not let it proceed against a schema that just changed underneath it.
"""
mock_store = MagicMock()
mock_store.has_pending_migration.return_value = True
mock_store.check_and_run_migrations.return_value = True
mocker.patch(
"paperless_ai.indexing.write_store",
return_value=mocker.MagicMock(
__enter__=mocker.MagicMock(return_value=mock_store),
__exit__=mocker.MagicMock(return_value=False),
),
)
mock_node = MagicMock()
mock_node.get_content.return_value = "fake node text"
mocker.patch(
"paperless_ai.indexing.build_document_node",
return_value=[mock_node],
)
doc = MagicMock(spec=Document)
doc.id = 1
indexing.llm_index_add_or_update_document(doc)
mock_store.upsert_document.assert_not_called()
def test_remove_document_uses_write_store(
self,
temp_llm_index_dir: Path,
mocker: pytest_mock.MockerFixture,
) -> None:
mock_store = MagicMock()
mock_store.has_pending_migration.return_value = False
mocker.patch(
"paperless_ai.indexing.write_store",
return_value=mocker.MagicMock(
@@ -777,6 +812,31 @@ class TestLlmIndexLocking:
mock_store.delete.assert_called_once_with("1")
def test_remove_document_skips_write_when_reembed_pending(
self,
temp_llm_index_dir: Path,
mocker: pytest_mock.MockerFixture,
) -> None:
"""A pending re-embed migration must block the delete too, for the
same consistency reason as the incremental-update path.
"""
mock_store = MagicMock()
mock_store.has_pending_migration.return_value = True
mock_store.check_and_run_migrations.return_value = True
mocker.patch(
"paperless_ai.indexing.write_store",
return_value=mocker.MagicMock(
__enter__=mocker.MagicMock(return_value=mock_store),
__exit__=mocker.MagicMock(return_value=False),
),
)
doc = MagicMock(spec=Document)
doc.id = 1
indexing.llm_index_remove_document(doc)
mock_store.delete.assert_not_called()
def test_update_llm_index_rebuild_uses_write_store(
self,
temp_llm_index_dir: Path,
@@ -849,6 +909,76 @@ class TestVectorStoreIndexing:
assert rows >= 1
class TestLlmIndexMigrate:
def test_noop_when_ai_disabled(self, mocker: pytest_mock.MockerFixture) -> None:
"""
GIVEN:
- AI/LLM index support is disabled in configuration
WHEN:
- llm_index_migrate() is called
THEN:
- No store is opened and no migration check runs
"""
mocker.patch(
"paperless_ai.indexing.AIConfig",
return_value=mocker.Mock(llm_index_enabled=False),
)
write_store_mock = mocker.patch("paperless_ai.indexing.write_store")
indexing.llm_index_migrate()
write_store_mock.assert_not_called()
def test_runs_pending_migration_when_enabled(
self,
mocker: pytest_mock.MockerFixture,
) -> None:
"""
GIVEN:
- AI/LLM index support is enabled
WHEN:
- llm_index_migrate() is called
THEN:
- The store is opened for write and a migration check runs
"""
mocker.patch(
"paperless_ai.indexing.AIConfig",
return_value=mocker.Mock(llm_index_enabled=True),
)
store_mock = mocker.MagicMock()
store_mock.has_pending_migration.return_value = False
write_store_cm = mocker.patch("paperless_ai.indexing.write_store")
write_store_cm.return_value.__enter__.return_value = store_mock
indexing.llm_index_migrate()
store_mock.has_pending_migration.assert_called_once()
def test_logs_warning_when_reembed_needed(
self,
mocker: pytest_mock.MockerFixture,
caplog: pytest.LogCaptureFixture,
) -> None:
"""
GIVEN:
- AI/LLM index support is enabled
- A pending migration requires re-embedding
WHEN:
- llm_index_migrate() is called
THEN:
- A warning directs the operator to run a manual rebuild, since
this automatic check must never re-embed on its own
"""
mocker.patch(
"paperless_ai.indexing.AIConfig",
return_value=mocker.Mock(llm_index_enabled=True),
)
store_mock = mocker.MagicMock()
store_mock.has_pending_migration.return_value = True
store_mock.check_and_run_migrations.return_value = True
write_store_cm = mocker.patch("paperless_ai.indexing.write_store")
write_store_cm.return_value.__enter__.return_value = store_mock
with caplog.at_level(logging.WARNING, logger="paperless_ai.indexing"):
indexing.llm_index_migrate()
assert "requires re-embedding" in caplog.text
@pytest.mark.django_db
class TestQuerySimilarDocuments:
def test_query_similar_documents_respects_allowed_ids(
+49 -2
View File
@@ -9,11 +9,11 @@ from llama_index.core.vector_stores.types import MetadataFilter
from llama_index.core.vector_stores.types import MetadataFilters
from llama_index.core.vector_stores.types import VectorStoreQuery
from paperless_ai.migrations import MIGRATIONS
from paperless_ai.migrations import Migration
from paperless_ai.vector_store import DB_FILENAME
from paperless_ai.vector_store import DEFAULT_TABLE_NAME
from paperless_ai.vector_store import MIGRATIONS
from paperless_ai.vector_store import SCHEMA_VERSION
from paperless_ai.vector_store import Migration
from paperless_ai.vector_store import PaperlessSqliteVecVectorStore
from paperless_ai.vector_store import _build_where
@@ -646,3 +646,50 @@ class TestMigrations:
assert result is True
assert self._schema_version(store) == 2
def test_has_pending_migration_false_when_no_table(
self,
store: PaperlessSqliteVecVectorStore,
) -> None:
"""
GIVEN:
- A vector store with no table created yet
WHEN:
- has_pending_migration() is checked
THEN:
- False is returned (nothing to migrate before anything exists)
"""
assert store.has_pending_migration() is False
def test_has_pending_migration_false_at_current_version(
self,
store: PaperlessSqliteVecVectorStore,
) -> None:
"""
GIVEN:
- A store at the current SCHEMA_VERSION
WHEN:
- has_pending_migration() is checked
THEN:
- False is returned
"""
store.add([make_node("a1", "1")])
assert store.has_pending_migration() is False
def test_has_pending_migration_true_when_behind(
self,
store: PaperlessSqliteVecVectorStore,
) -> None:
"""
GIVEN:
- A store whose schema_version has been forced behind SCHEMA_VERSION
WHEN:
- has_pending_migration() is checked
THEN:
- True is returned
"""
store.add([make_node("a1", "1")])
store.client.execute(
"UPDATE index_meta SET value = '0' WHERE key = 'schema_version'",
)
assert store.has_pending_migration() is True
+35 -45
View File
@@ -2,16 +2,12 @@ import json
import logging
import sqlite3
import struct
from collections.abc import Callable
from collections.abc import Iterator
from collections.abc import Sequence
from contextlib import contextmanager
from dataclasses import dataclass
from dataclasses import field
from pathlib import Path
from types import TracebackType
from typing import Any
from typing import Literal
import sqlite_vec
from llama_index.core.bridge.pydantic import PrivateAttr
@@ -26,6 +22,9 @@ from llama_index.core.vector_stores.types import VectorStoreQueryResult
from llama_index.core.vector_stores.utils import metadata_dict_to_node
from llama_index.core.vector_stores.utils import node_to_metadata_dict
from paperless_ai.migrations import MIGRATIONS
from paperless_ai.migrations import Migration
logger = logging.getLogger("paperless_ai.vector_store")
DB_FILENAME = "llmindex.db"
@@ -53,38 +52,6 @@ COMPACT_BATCH_SIZE = 500
_FILTER_COLUMNS = frozenset({"document_id", "modified"})
@dataclass
class Migration:
"""A schema migration for the sqlite-vec vector store.
kind="structural": rows are copied into a new-schema file with no
re-embedding needed. Supply ``apply(src_conn, dst_conn, dim)`` which
must create the vec0 table in ``dst_conn``, copy all rows from
``src_conn``, and write ``dim`` / ``embed_model`` / ``total_inserts`` to
``dst_conn``'s ``index_meta``. ``schema_version`` is written by the
migration runner after ``apply`` returns.
kind="re-embed": the new schema requires fresh embeddings.
``check_and_run_migrations()`` returns True when it encounters one of
these so the caller can force a full rebuild (which recreates the table
at the current SCHEMA_VERSION).
"""
from_version: int
to_version: int
kind: Literal["structural", "re-embed"]
description: str
apply: Callable[[sqlite3.Connection, sqlite3.Connection, int], None] | None = field(
default=None,
repr=False,
)
# Registry of all schema migrations in order. Empty at v1 -- this is the
# baseline. Add entries here (and bump SCHEMA_VERSION) when the schema changes.
MIGRATIONS: list[Migration] = []
def _pack(embedding: Sequence[float]) -> bytes:
return struct.pack(f"{len(embedding)}f", *embedding)
@@ -551,6 +518,31 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
Path(compact_path).replace(db_path)
self._conn = self._open_connection(db_path)
def _stored_schema_version(self) -> int | None:
"""The schema_version recorded in index_meta, or None if no table
exists. A missing key (a store predating version tracking) is
treated as SCHEMA_VERSION -- i.e. already current -- since no
migration in MIGRATIONS targets a version before tracking began.
"""
if not self.table_exists():
return None
raw = self._meta_get("schema_version")
return int(raw) if raw is not None else SCHEMA_VERSION
def has_pending_migration(self) -> bool:
"""Cheaply check whether a migration is pending, with no exclusive
access needed -- just a metadata read under the connection callers
already hold via the write FileLock.
Callers should only pay for check_and_run_migrations()'s exclusive
access (a structural migration's file swap must not run while
readers are active) when this returns True, so that the common
case -- already at SCHEMA_VERSION -- never contends with readers
or a concurrent compaction.
"""
current = self._stored_schema_version()
return current is not None and current < SCHEMA_VERSION
def check_and_run_migrations(self) -> bool:
"""Apply any pending schema migrations to the store.
@@ -559,15 +551,13 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
this method returns True when one is encountered so the caller can
force a full rebuild (which recreates the table at SCHEMA_VERSION).
Must be called under the write FileLock. No-op when the table does
not exist or is already at SCHEMA_VERSION.
Must be called under the write FileLock, with readers excluded (see
has_pending_migration() for a cheap pre-check that avoids paying for
that exclusion in the common case). No-op when the table does not
exist or is already at SCHEMA_VERSION.
"""
if not self.table_exists():
return False
raw = self._meta_get("schema_version")
current = int(raw) if raw is not None else SCHEMA_VERSION
if current >= SCHEMA_VERSION:
current = self._stored_schema_version()
if current is None or current >= SCHEMA_VERSION:
return False
pending = sorted(
@@ -579,7 +569,7 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
if migration.kind == "re-embed":
logger.warning(
"LLM index schema v%d -> v%d requires re-embedding (%s); "
"forcing full rebuild.",
"the caller must force a rebuild.",
migration.from_version,
migration.to_version,
migration.description,