Compare commits

...
Author SHA1 Message Date
stumpylogandClaude Sonnet 5 48b73562f0 Perf/fix: stop building a full-library IN filter for unrestricted RAG context
get_context_for_document() always materialized every document id a user
can see into a Python list, even for a superuser (or no user at all),
passing it through as a SQL IN filter. For a superuser, that's the whole
library:

- Past ~32,763 documents, this crashes outright:
  sqlite3.OperationalError: too many SQL variables (SQLite's
  SQLITE_MAX_VARIABLE_NUMBER is 32766 by default, and the query already
  binds embedding + k + a NE self-exclusion clause alongside the ids).
- Below that cliff, vec0's IN-list evaluation is a nested loop (strncmp
  per row per allowed id), so it's quadratic in library size for no
  reason -- the filter was never going to exclude anything.

get_objects_for_user_owner_aware() already returns every Document for a
superuser (guardian's own with_superuser shortcut), so skipping straight
to document_ids=None changes nothing about which documents are
considered, only how we get there.

Also:
- Drop the pointless sorted() in _document_id_filters(): a SQL IN clause
  doesn't care about order, so it was pure overhead on every call.
- Add a hard guard in _build_where() so a future regression (or a
  legitimately huge permission-restricted user) fails closed -- no rows,
  a logged warning -- instead of a cryptic OperationalError. Since this
  filter scopes document access, failing closed rather than skipping the
  filter is the only safe way to handle an oversized list.

First item from VECTOR_STORE_PERF_BACKLOG.md (an audit done alongside
perf/13314-vecstore-point-delete, deferred to its own branch since that
one was already large).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 13:00:48 -07:00
stumpylogandClaude Sonnet 5 15938945d7 Simplify: dedupe file-swap rebuild, migration-check locking, and test mocks
Prompted by an automated simplification review of the branch. Applies the
highest-value findings, all behavior-preserving (full suite still green):

- vector_store.py: extract _rebuild_file() (temp-file lifecycle: open,
  populate, swap-in-or-discard-on-failure) and _rebuild_into() (create
  table, copy meta, stream rows) out of compact() and the v1->v2
  migration, which had duplicated both almost verbatim. The migration
  module shrinks to a single _rebuild_into() call instead of reaching
  into four PaperlessSqliteVecVectorStore privates.
- vector_store.py: extract _stored_schema_version() out of
  has_pending_migration() and check_and_run_migrations(), which
  duplicated the same schema_version read and its "missing key means
  current" default.
- indexing.py: extract _with_exclusive_access(), collapsing three
  identical _exclude_readers()/Timeout/log-and-skip blocks (compaction
  in update_llm_index() and llm_index_compact(), the migration check in
  _check_and_run_migrations()) to one line each. Also trims
  _check_and_run_migrations()'s docstring, which had grown to restate
  content already documented on has_pending_migration() and
  check_and_run_migrations(), including a claim that was now stale
  (llm_index_migrate() is a second consumer of the re-embed signal).
- test_ai_indexing.py: add a mock_store fixture for the
  write_store()-yields-a-MagicMock pattern that 8 tests were hand-rolling
  identically.
- migrations/__init__.py: consolidate the "how to add a migration"
  instructions to one place instead of three.
- docs/administration.md: fix two now-inaccurate claims -- the bare-metal
  migrate step doesn't run "automatically" (it's the manual step being
  documented), and the self-contradictory "if enabled... no-op if
  disabled" phrasing on the same step.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 12:32:59 -07:00
stumpylogandClaude Sonnet 5 470bcc1751 Refactor: address PR review comments
- Move schema migrations into their own paperless_ai/migrations package,
  one module per migration (mNNNN_description.py -- a leading digit isn't
  a valid Python identifier, unlike Django's own numbered migrations,
  which load via importlib.import_module() rather than a static import
  statement), establishing the pattern before more migrations accumulate
  in vector_store.py.
- Drop the redundant "if the LLM index is enabled..." clause from the
  Docker upgrade note in docs/administration.md -- it's already covered
  by the LLM index section below, and calling out just one of several
  auto-applied migrations there was incomplete.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 11:45:18 -07:00
stumpylogandClaude Sonnet 5 cd7038a7a4 Fix: don't gate document_ids-scoped LLM index updates on Document.modified
bulk_edit.py's tag/correspondent/document_type/storage_path/custom_field
helpers write via queryset.update() or direct M2M/through-model bulk
operations, none of which call Document.save() -- so Document.modified's
auto_now never fires. Comparing against get_modified_times() for a
document_ids-scoped update_llm_index() call was silently skipping reindex
of documents whose embedded metadata had just changed via a bulk edit.
The modified-time comparison is now only used for the unscoped,
full-library incremental scan, where it is still needed.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 11:13:05 -07:00
stumpylogandClaude Sonnet 5 df8a6cc715 Feature: run pending LLM index migrations automatically on startup
Add document_llmindex migrate, a cheap check-only path (no reindex) safe
to run unconditionally: has_pending_migration() short-circuits to a
metadata-only read once the store is current, so a healthy install pays
almost nothing. If a pending migration would require re-embedding, it
only logs a warning and leaves the index as-is -- re-embedding can be
slow and, for a metered embedding backend, cost money, so it stays a
deliberate manual action (document_llmindex rebuild), never automatic.

Wire it into the Docker image as a new init-llmindex-migrate oneshot
(modeled on init-search-index, gated on init-migrations), and document
the equivalent manual step for bare-metal upgrades.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 10:48:24 -07:00
stumpylogandClaude Sonnet 5 b21d50645d Fix: run pending structural migrations before delete()/upsert_document()
Only update_llm_index()'s nightly rebuild ran check_and_run_migrations(),
but llm_index_remove_document()/llm_index_add_or_update_document() call
delete()/upsert_document() directly and immediately (on every document
delete/edit). On upgrade, every pre-existing document is "pre-migration"
until the next scheduled rebuild -- up to 24h by default -- so a delete or
edit in that window would find zero document_chunks rows and silently
leave stale/orphaned vec0 rows behind.

Add has_pending_migration(), a cheap metadata-only read that needs no
exclusive access, so normal writes only pay for check_and_run_migrations()'s
exclusive locking when a migration is actually pending -- otherwise they'd
be gated by the same lock compaction uses, which
test_normal_write_is_not_gated_by_the_compaction_lock exists to forbid.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 10:30:04 -07:00
stumpylog 24b86770df Perf: store document_chunks.document_id as INTEGER
document_chunks is a plain SQLite table, not a vec0 virtual table, so
(unlike vec0's own TEXT metadata column) it gets normal type-affinity
coercion between the TEXT document ids used elsewhere in this module and
an INTEGER column -- verified empirically. INTEGER keys make the
per-document delete lookup this table exists for cheaper: smaller index
entries and integer rather than text B-tree comparisons.
2026-07-28 09:39:02 -07:00
stumpylog 6a2972f313 Test: add missing document_chunks coverage and GIVEN/WHEN/THEN docstrings
Cover compact()'s handling of a churned document's final chunk generation,
delete() on a never-indexed document, upsert_document() clearing to an empty
node list, and the v1->v2 migration's embed_model preservation and
idempotency, per repo convention.
2026-07-28 09:36:58 -07:00
Trenton HolmesandClaude Sonnet 5 fa9741b555 Simplify: dedupe row-copy/migration-test helpers, type fixture params
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 08:56:16 -07:00
Trenton HolmesandClaude Sonnet 5 63c9430eec Perf: point-delete vector store chunks by id instead of a document_id scan
vec0 only gets an efficient lookup on a metadata column inside a KNN
query; a plain document_id-filtered DELETE is a full table scan
regardless of index size. Add a document_chunks side table (plain
SQLite, real index) to look up chunk ids for a document, then delete
each by its id primary key instead. Includes a v1 -> v2 migration to
backfill document_chunks for indexes created before this existed.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 08:56:16 -07:00
16 changed files with 1741 additions and 380 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 LLM index schema..."
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
+20 -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. Apply any pending LLM index schema migrations.
```shell-session
cd src
python3 manage.py document_llmindex migrate
```
This is a no-op if the index is already up to date, or if the LLM index
is disabled, 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
@@ -544,6 +554,15 @@ scheduled task runs.
Specify `compact` to reclaim space and optimize the on-disk vector store.
Specify `migrate` to apply any pending index schema migrations without a full reindex.
This is a no-op if the index is already up to date, so it is safe to run on every
startup or upgrade; the container's startup sequence runs it automatically, and the
[bare-metal upgrade steps](#bare-metal-updating) include it as a manual step. If a
pending migration would require re-embedding every document, `migrate` only logs a
warning and leaves the index as-is -- re-embedding can be slow and, for a metered
embedding backend, cost money, so it is never triggered automatically. Run `rebuild`
yourself when you are ready.
!!! note
These commands have no effect unless AI is enabled and an embedding backend is
@@ -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(
@@ -8,6 +8,7 @@ if TYPE_CHECKING:
from pytest_mock import MockerFixture
_COMPACT = "documents.management.commands.document_llmindex.llm_index_compact"
_MIGRATE = "documents.management.commands.document_llmindex.llm_index_migrate"
_INDEX = "documents.management.commands.document_llmindex.llmindex_index"
@@ -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,
+19 -12
View File
@@ -96,19 +96,26 @@ def get_context_for_document(
user: User | None = None,
max_docs: int = 5,
) -> str:
visible_documents = (
get_objects_for_user_owner_aware(
user,
"view_document",
Document,
)
if user
else None
)
# None means "no restriction" to query_similar_documents. A superuser
# (like no user at all) can see every document, so skip materializing
# every visible pk into a Python list and passing it through as a SQL
# IN filter: for a large library that is a wasted quadratic scan in the
# vector store at best, and past ~32,763 documents a hard
# sqlite3.OperationalError (SQLite's bound-parameter limit) at worst.
# get_objects_for_user_owner_aware() would return every Document for a
# superuser anyway (guardian's own with_superuser shortcut), so this
# changes nothing about which documents are considered -- only how we
# get there.
visible_document_ids = (
list(visible_documents.values_list("pk", flat=True))
if visible_documents is not None
else None
None
if user is None or user.is_superuser
else list(
get_objects_for_user_owner_aware(
user,
"view_document",
Document,
).values_list("pk", flat=True),
)
)
similar_docs = query_similar_documents(
document=doc,
+79 -29
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.
@@ -301,7 +319,7 @@ def _document_id_filters(doc_ids):
MetadataFilter(
key="document_id",
operator=FilterOperator.IN,
value=sorted(doc_ids),
value=list(doc_ids),
),
],
)
@@ -324,6 +342,21 @@ def _exclude_document_id_filter(document_id: int | str):
)
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 update_llm_index(
*,
iter_wrapper: IterWrapper[Document] = identity,
@@ -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.",
@@ -396,12 +421,24 @@ def update_llm_index(
if document_ids is not None
else documents
)
existing = store.get_modified_times()
# When document_ids is given, the caller already knows exactly
# which documents to reindex (e.g. a bulk edit) -- trust it and
# skip the modified-time comparison entirely. Bulk edits (tags,
# correspondent, document type, storage path, custom fields)
# write via queryset.update()/M2M bulk operations, which bypass
# Document.modified's auto_now, so comparing against
# get_modified_times() here would silently skip reindexing
# documents whose embedded metadata just changed. The comparison
# is only meaningful for the unscoped, full-library scan, where
# it avoids re-embedding documents that have not changed.
existing = store.get_modified_times() if document_ids is None else None
changed = 0
for document in iter_wrapper(scoped_documents):
doc_id = str(document.id)
if existing.get(doc_id) == document.modified.isoformat():
continue
if existing is not None:
stored_modified = existing.get(doc_id)
if stored_modified == document.modified.isoformat():
continue
nodes = build_document_node(document, chunk_size=chunk_size)
_embed_nodes(nodes, embed_model)
store.upsert_document(doc_id, nodes)
@@ -412,14 +449,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 +464,45 @@ 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:
_check_and_run_migrations(store)
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:
_check_and_run_migrations(store)
store.delete(str(document.id))
+58
View File
@@ -0,0 +1,58 @@
"""Schema migrations for the sqlite-vec vector store.
Each migration lives in its own module here, named ``mNNNN_description.py``
(e.g. ``m0001_add_document_chunks.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()`` (most likely just a call to
``PaperlessSqliteVecVectorStore._rebuild_into()``, see ``Migration`` below),
and appends a ``Migration`` to ``MIGRATIONS``; then import that module at the
bottom of ``vector_store.py`` and bump ``SCHEMA_VERSION`` there.
"""
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 the vec0 table in ``dst_conn`` and copy ``src_conn``'s rows
and index_meta into it -- usually just a call to
``PaperlessSqliteVecVectorStore._rebuild_into(src_conn, dst_conn, dim)``.
``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] = []
@@ -0,0 +1,41 @@
import sqlite3
from paperless_ai.migrations import MIGRATIONS
from paperless_ai.migrations import Migration
from paperless_ai.vector_store import PaperlessSqliteVecVectorStore
def _migrate_v1_to_v2_add_document_chunks(
src_conn: sqlite3.Connection,
dst_conn: sqlite3.Connection,
dim: int,
) -> None:
"""v1 -> v2: backfill the document_chunks side table.
document_chunks (see PaperlessSqliteVecVectorStore._open_connection) lets
delete()/upsert_document() find a document's chunk ids without a vec0
full table scan on the document_id metadata column. Every row written
before this migration predates that table, so without backfilling,
deleting a pre-migration document would find zero chunk ids and leave its
vec0 rows permanently orphaned. Backfilling is just the same rebuild
compact() uses (_rebuild_into), minus "schema_version" -- the caller
(_run_structural_migration) sets that to this migration's target version
instead of preserving the source's.
"""
PaperlessSqliteVecVectorStore._rebuild_into(
src_conn,
dst_conn,
dim,
meta_keys=("dim", "embed_model"),
)
MIGRATIONS.append(
Migration(
from_version=1,
to_version=2,
kind="structural",
description="add document_chunks side table for O(1) per-document deletes",
apply=_migrate_v1_to_v2_add_document_chunks,
),
)
@@ -3,6 +3,8 @@ 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 documents.models import Document
@@ -278,3 +280,104 @@ def test_get_context_for_document_no_similar_docs(mock_document):
with patch("paperless_ai.ai_classifier.query_similar_documents", return_value=[]):
result = get_context_for_document(mock_document)
assert result == ""
class TestGetContextForDocumentVisibility:
"""get_context_for_document must not materialize every visible document
id for a user who can already see the whole library: a superuser (like
no user at all) gets document_ids=None (no restriction) straight
through to query_similar_documents(), instead of a full-library IN
filter that is wasteful at best and, past ~32,763 documents, a hard
sqlite3.OperationalError at worst (SQLite's bound-parameter limit).
"""
def test_skips_permission_lookup_for_superuser(
self,
mock_document: MagicMock,
mock_similar_documents: list[MagicMock],
mocker: pytest_mock.MockerFixture,
) -> None:
"""
GIVEN:
- A superuser
WHEN:
- get_context_for_document() is called
THEN:
- get_objects_for_user_owner_aware() is never called, and
query_similar_documents() is called with document_ids=None
"""
mock_query = mocker.patch(
"paperless_ai.ai_classifier.query_similar_documents",
return_value=mock_similar_documents,
)
mock_get_objects = mocker.patch(
"paperless_ai.ai_classifier.get_objects_for_user_owner_aware",
)
user = mocker.MagicMock(spec=User)
user.is_superuser = True
get_context_for_document(mock_document, user, max_docs=2)
mock_get_objects.assert_not_called()
assert mock_query.call_args.kwargs["document_ids"] is None
def test_skips_permission_lookup_when_no_user(
self,
mock_document: MagicMock,
mock_similar_documents: list[MagicMock],
mocker: pytest_mock.MockerFixture,
) -> None:
"""
GIVEN:
- No user (user=None)
WHEN:
- get_context_for_document() is called
THEN:
- get_objects_for_user_owner_aware() is never called, and
query_similar_documents() is called with document_ids=None
"""
mock_query = mocker.patch(
"paperless_ai.ai_classifier.query_similar_documents",
return_value=mock_similar_documents,
)
mock_get_objects = mocker.patch(
"paperless_ai.ai_classifier.get_objects_for_user_owner_aware",
)
get_context_for_document(mock_document, None, max_docs=2)
mock_get_objects.assert_not_called()
assert mock_query.call_args.kwargs["document_ids"] is None
def test_restricts_to_visible_documents_for_non_superuser(
self,
mock_document: MagicMock,
mock_similar_documents: list[MagicMock],
mocker: pytest_mock.MockerFixture,
) -> None:
"""
GIVEN:
- A non-superuser with a specific set of visible documents
WHEN:
- get_context_for_document() is called
THEN:
- query_similar_documents() is called with exactly that user's
visible document ids, unchanged from before this optimization
"""
mock_query = mocker.patch(
"paperless_ai.ai_classifier.query_similar_documents",
return_value=mock_similar_documents,
)
mock_queryset = mocker.MagicMock()
mock_queryset.values_list.return_value = [1, 2, 3]
mock_get_objects = mocker.patch(
"paperless_ai.ai_classifier.get_objects_for_user_owner_aware",
return_value=mock_queryset,
)
user = mocker.MagicMock(spec=User)
user.is_superuser = False
get_context_for_document(mock_document, user, max_docs=2)
mock_get_objects.assert_called_once_with(user, "view_document", Document)
assert mock_query.call_args.kwargs["document_ids"] == [1, 2, 3]
+233 -36
View File
@@ -21,6 +21,7 @@ from documents.signals import document_consumption_finished
from documents.signals import document_updated
from documents.tests.factories import DocumentFactory
from documents.tests.factories import PaperlessTaskFactory
from documents.tests.factories import TagFactory
from paperless.models import ApplicationConfiguration
from paperless_ai import indexing
from paperless_ai.tests.conftest import FakeEmbedding
@@ -36,6 +37,23 @@ def real_document(db: None) -> Document:
)
@pytest.fixture
def mock_store(mocker: pytest_mock.MockerFixture) -> MagicMock:
"""The MagicMock store yielded by every ``with write_store() as store:``
block, for tests that only care what indexing.py does with the store,
not what the store itself does.
"""
store = mocker.MagicMock()
mocker.patch(
"paperless_ai.indexing.write_store",
return_value=mocker.MagicMock(
__enter__=mocker.MagicMock(return_value=store),
__exit__=mocker.MagicMock(return_value=False),
),
)
return store
@pytest.mark.django_db
def test_build_document_node(real_document: Document) -> None:
nodes = indexing.build_document_node(real_document)
@@ -347,6 +365,62 @@ def test_update_llm_index_partial_update(
assert after[str(doc2.pk)] == before[str(doc2.pk)]
class TestUpdateLlmIndexScopedDocumentIds:
"""A document_ids-scoped update must trust the caller and reindex every scoped
document, never gating on Document.modified.
bulk_edit.py's add_tag/remove_tag/modify_tags/set_correspondent/set_document_type/
set_storage_path/modify_custom_fields all write via queryset.update() or direct
M2M/through-model bulk operations -- none of which call Document.save(), so
Document.modified's auto_now never fires. Comparing against
get_modified_times() for a document_ids-scoped call would therefore skip
reindexing documents whose embedded tags/correspondent/etc. just changed.
"""
@pytest.mark.django_db
def test_scoped_update_reindexes_despite_unchanged_modified(
self,
temp_llm_index_dir: Path,
mock_embed_model: FakeEmbedding,
) -> None:
"""A document_ids-scoped update must pick up a tag added via the M2M
manager's bulk_create path, even though that leaves modified untouched.
Steps:
1. Build an initial index for a document with no tags.
2. Add a tag via a direct through-model bulk_create, mirroring
bulk_edit.add_tag -- this does not call Document.save().
3. Call update_llm_index(rebuild=False, document_ids=[doc.pk]).
4. Assert the stored node metadata now includes the new tag.
"""
# Step 1
tag = TagFactory.create(name="Important")
doc = DocumentFactory.create(title="Test Document", added=timezone.now())
indexing.update_llm_index(rebuild=True)
modified_before = doc.modified
# Step 2: bulk-add the tag the way bulk_edit.add_tag does -- a direct
# through-model insert, no Document.save().
DocumentTagRelationship = Document.tags.through
DocumentTagRelationship.objects.bulk_create(
[DocumentTagRelationship(document_id=doc.pk, tag_id=tag.pk)],
)
doc.refresh_from_db()
assert doc.modified == modified_before, (
"Precondition failed: expected modified to be unchanged after a "
"through-model bulk tag add"
)
# Step 3
result = indexing.update_llm_index(rebuild=False, document_ids=[doc.pk])
assert result == "LLM index updated successfully."
# Step 4
with indexing.get_vector_store() as store:
nodes = store.get_nodes()
assert any(tag.name in node.metadata.get("tags", []) for node in nodes)
@pytest.mark.django_db
def test_add_or_update_document_updates_existing_entry(
temp_llm_index_dir: Path,
@@ -705,23 +779,104 @@ class TestLlmIndexAddOrUpdateDocumentEmptyContent:
@pytest.mark.django_db
def test_llm_index_compact_uses_force(
temp_llm_index_dir: Path,
mocker: pytest_mock.MockerFixture,
mock_store: MagicMock,
) -> None:
"""compact must use force=True to rebuild the table and reclaim space immediately."""
mock_store = mocker.MagicMock()
mocker.patch(
"paperless_ai.indexing.write_store",
return_value=mocker.MagicMock(
__enter__=mocker.MagicMock(return_value=mock_store),
__exit__=mocker.MagicMock(return_value=False),
),
)
indexing.llm_index_compact()
mock_store.compact.assert_called_once_with(force=True)
@pytest.mark.django_db
class TestLlmIndexMigrate:
"""llm_index_migrate() is the cheap, startup-safe migration check -- see
the init-llmindex-migrate container step and the bare-metal upgrade docs.
"""
def test_skips_when_llm_index_disabled(
self,
temp_llm_index_dir: Path,
mocker: pytest_mock.MockerFixture,
) -> None:
"""
GIVEN:
- The LLM index is disabled
WHEN:
- llm_index_migrate() is called
THEN:
- The store is never opened (no stray db file for users who
never enabled AI features)
"""
mock_config = mocker.MagicMock()
mock_config.llm_index_enabled = False
mocker.patch("paperless_ai.indexing.AIConfig", return_value=mock_config)
mock_write_store = mocker.patch("paperless_ai.indexing.write_store")
indexing.llm_index_migrate()
mock_write_store.assert_not_called()
def test_runs_pending_structural_migration(
self,
temp_llm_index_dir: Path,
mocker: pytest_mock.MockerFixture,
mock_store: MagicMock,
caplog: pytest.LogCaptureFixture,
) -> None:
"""
GIVEN:
- The LLM index is enabled and a structural migration is pending
WHEN:
- llm_index_migrate() is called
THEN:
- check_and_run_migrations() runs, and no re-embed warning is
logged (the pending migration was structural, not re-embed)
"""
mock_config = mocker.MagicMock()
mock_config.llm_index_enabled = True
mocker.patch("paperless_ai.indexing.AIConfig", return_value=mock_config)
mock_store.has_pending_migration.return_value = True
mock_store.check_and_run_migrations.return_value = False
with caplog.at_level("WARNING"):
indexing.llm_index_migrate()
mock_store.check_and_run_migrations.assert_called_once()
assert "re-embedding" not in caplog.text
def test_warns_without_rebuilding_when_reembed_pending(
self,
temp_llm_index_dir: Path,
mocker: pytest_mock.MockerFixture,
mock_store: MagicMock,
caplog: pytest.LogCaptureFixture,
) -> None:
"""
GIVEN:
- The LLM index is enabled and a pending migration requires
re-embedding
WHEN:
- llm_index_migrate() is called
THEN:
- A warning is logged telling the admin to rebuild manually, but
no rebuild is triggered automatically -- re-embedding can be
slow and, for a metered embedding backend, cost money, so it
must be a deliberate user action, never an automatic one
"""
mock_config = mocker.MagicMock()
mock_config.llm_index_enabled = True
mocker.patch("paperless_ai.indexing.AIConfig", return_value=mock_config)
mock_store.has_pending_migration.return_value = True
mock_store.check_and_run_migrations.return_value = True
with caplog.at_level("WARNING"):
indexing.llm_index_migrate()
assert "re-embedding" in caplog.text
mock_store.drop_table.assert_not_called()
mock_store.add.assert_not_called()
@pytest.mark.django_db
class TestLlmIndexLocking:
"""Index mutation functions must go through write_store(), which holds the lock.
@@ -734,16 +889,9 @@ class TestLlmIndexLocking:
self,
temp_llm_index_dir: Path,
mock_embed_model: FakeEmbedding,
mock_store: MagicMock,
mocker: pytest_mock.MockerFixture,
) -> None:
mock_store = MagicMock()
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(
@@ -757,40 +905,89 @@ class TestLlmIndexLocking:
mock_store.upsert_document.assert_called_once()
@pytest.mark.parametrize("has_pending", [True, False])
def test_add_or_update_document_runs_migration_check_only_when_pending(
self,
temp_llm_index_dir: Path,
mock_embed_model: FakeEmbedding,
mock_store: MagicMock,
mocker: pytest_mock.MockerFixture,
*,
has_pending: bool,
) -> None:
"""
GIVEN:
- A document to add/update, and a store reporting whether a
migration is pending
WHEN:
- llm_index_add_or_update_document() is called
THEN:
- check_and_run_migrations() runs only when has_pending_migration()
is True, so a normal upsert never pays for the exclusive access
that check_and_run_migrations() requires
- upsert_document() is called either way
"""
mock_store.has_pending_migration.return_value = has_pending
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 = DocumentFactory.build(id=1)
indexing.llm_index_add_or_update_document(doc)
assert mock_store.check_and_run_migrations.called is has_pending
mock_store.upsert_document.assert_called_once()
def test_remove_document_uses_write_store(
self,
temp_llm_index_dir: Path,
mocker: pytest_mock.MockerFixture,
mock_store: MagicMock,
) -> None:
mock_store = MagicMock()
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_called_once_with("1")
@pytest.mark.parametrize("has_pending", [True, False])
def test_remove_document_runs_migration_check_only_when_pending(
self,
temp_llm_index_dir: Path,
mock_store: MagicMock,
*,
has_pending: bool,
) -> None:
"""
GIVEN:
- A document to remove, and a store reporting whether a
migration is pending
WHEN:
- llm_index_remove_document() is called
THEN:
- check_and_run_migrations() runs only when has_pending_migration()
is True, so a normal delete never pays for the exclusive access
that check_and_run_migrations() requires (see
test_normal_write_is_not_gated_by_the_compaction_lock)
- delete() is called either way
"""
mock_store.has_pending_migration.return_value = has_pending
doc = DocumentFactory.build(id=1)
indexing.llm_index_remove_document(doc)
assert mock_store.check_and_run_migrations.called is has_pending
mock_store.delete.assert_called_once_with("1")
def test_update_llm_index_rebuild_uses_write_store(
self,
temp_llm_index_dir: Path,
mock_embed_model: FakeEmbedding,
mock_store: MagicMock,
mocker: pytest_mock.MockerFixture,
) -> None:
mock_store = MagicMock()
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_qs = MagicMock()
mock_qs.exists.return_value = True
mock_qs.__iter__ = MagicMock(return_value=iter([]))
File diff suppressed because it is too large Load Diff
+240 -120
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,15 +22,27 @@ 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"
DEFAULT_TABLE_NAME = "documents"
# Current schema version. Written to index_meta at table creation and bumped
# whenever a Migration is added to MIGRATIONS. check_and_run_migrations() uses
# this to decide which migrations to run on an existing store.
SCHEMA_VERSION = 1
_INSERT = (
"INSERT INTO "
+ DEFAULT_TABLE_NAME
+ " (id, document_id, modified, node_content, embedding) VALUES (?, ?, ?, ?, ?)"
)
_INSERT_CHUNK_INDEX = (
"INSERT INTO document_chunks (chunk_id, document_id) VALUES (?, ?)"
)
# Current schema version. Bump when adding a migration -- see
# paperless_ai/migrations/__init__.py for the full procedure.
SCHEMA_VERSION = 2
# compact(): rebuild when the cumulative rowid count exceeds this multiple of
# the live row count. DELETEs on vec0 tables never reclaim space (upstream
@@ -52,37 +60,15 @@ COMPACT_BATCH_SIZE = 500
# construction.
_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] = []
# _build_where(): the largest IN value list translated into bound SQL
# parameters. SQLite's own hard limit (SQLITE_MAX_VARIABLE_NUMBER) is 32766
# by default; this leaves headroom below that for the query's other bound
# parameters (the embedding blob, k, and any NE clause) and for the limit
# itself to move. An IN filter this large should not happen in practice --
# callers are expected to pass None (no filter) rather than every id when
# the filter would not actually narrow anything -- so this is a guard
# against a future regression, not a normal code path.
_MAX_IN_VALUES = 32700
def _pack(embedding: Sequence[float]) -> bytes:
@@ -93,6 +79,42 @@ def _unpack(blob: bytes) -> list[float]:
return list(struct.unpack(f"{len(blob) // 4}f", blob))
def _copy_rows(src_conn: sqlite3.Connection, dst_conn: sqlite3.Connection) -> int:
"""Copy every live vec0 row from ``src_conn`` into ``dst_conn``, recording
each one in ``dst_conn``'s document_chunks side table. Returns the number
of rows copied. The caller owns ``dst_conn``'s transaction.
Rows are streamed from the source cursor in batches instead of being
materialized all at once, so a large index does not cause an OOM during a
routine compaction or migration.
"""
src_cursor = src_conn.execute(
"SELECT id, document_id, modified, node_content, embedding FROM "
+ DEFAULT_TABLE_NAME,
)
copied = 0
while batch := src_cursor.fetchmany(COMPACT_BATCH_SIZE):
dst_conn.executemany(
_INSERT,
[
(
r["id"],
r["document_id"],
r["modified"],
r["node_content"],
bytes(r["embedding"]),
)
for r in batch
],
)
dst_conn.executemany(
_INSERT_CHUNK_INDEX,
[(r["id"], r["document_id"]) for r in batch],
)
copied += len(batch)
return copied
def _build_where(filters: MetadataFilters | None) -> tuple[str, list[str]]:
"""Translate the EQ / IN / NE filters we use into a parameterized SQL clause
on vec0 metadata columns. Returns ("", []) when there is nothing to filter.
@@ -113,6 +135,21 @@ def _build_where(filters: MetadataFilters | None) -> tuple[str, list[str]]:
if not values: # pragma: no cover
clauses.append("1 = 0")
continue
if len(values) > _MAX_IN_VALUES:
# Fail closed (see the empty-clauses case below) rather than
# let SQLite raise "too many SQL variables" past its own
# limit: this filter scopes document access, so an IN list
# too large to safely bind must match no rows, never widen
# the scope to "everything" by accident.
logger.warning(
"Refusing to build an IN filter on %r with %d values "
"(over the %d-value safety limit); returning no rows.",
f.key,
len(values),
_MAX_IN_VALUES,
)
clauses.append("1 = 0")
continue
placeholders = ",".join("?" for _ in values)
clauses.append(f"{f.key} IN ({placeholders})")
params.extend(values)
@@ -189,6 +226,24 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
conn.execute(
"CREATE TABLE IF NOT EXISTS index_meta (key TEXT PRIMARY KEY, value TEXT)",
)
# vec0 metadata columns only get an efficient lookup path inside a KNN
# (MATCH) query; a plain `WHERE document_id = ?` is a full table scan
# regardless of index size. This plain, indexed table is how delete()/
# upsert_document() find a document's chunk ids without that scan.
# document_id is INTEGER here (unlike vec0's own TEXT metadata
# column): this is a normal SQLite table, so standard type affinity
# correctly coerces the TEXT document ids written/looked-up
# elsewhere in this module -- it does not share vec0's own
# metadata-column comparison code, which silently mismatches
# non-TEXT bound values instead of coercing them.
conn.execute(
"CREATE TABLE IF NOT EXISTS document_chunks "
"(chunk_id TEXT PRIMARY KEY, document_id INTEGER NOT NULL)",
)
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_document_chunks_document_id "
"ON document_chunks (document_id)",
)
return conn
@property
@@ -223,13 +278,17 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
else:
self._conn.execute("COMMIT")
def _meta_get(self, key: str) -> str | None:
row = self._conn.execute(
@staticmethod
def _meta_get_on(conn: sqlite3.Connection, key: str) -> str | None:
row = conn.execute(
"SELECT value FROM index_meta WHERE key = ?",
(key,),
).fetchone()
return row["value"] if row else None
def _meta_get(self, key: str) -> str | None:
return self._meta_get_on(self._conn, key)
@staticmethod
def _meta_set_on(conn: sqlite3.Connection, key: str, value: str) -> None:
conn.execute(
@@ -259,6 +318,7 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
def drop_table(self) -> None:
self._conn.execute("DROP TABLE IF EXISTS " + DEFAULT_TABLE_NAME)
self._conn.execute("DELETE FROM index_meta")
self._conn.execute("DELETE FROM document_chunks")
def stored_model_name(self) -> str | None:
"""Return the embedding model name recorded at table creation, or None."""
@@ -325,11 +385,37 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
_pack(node.get_embedding()),
)
_INSERT = (
"INSERT INTO "
+ DEFAULT_TABLE_NAME
+ " (id, document_id, modified, node_content, embedding) VALUES (?, ?, ?, ?, ?)"
)
def _index_chunks(self, rows: list[tuple[str, str, str, str, bytes]]) -> None:
"""Record each row's (chunk_id, document_id) in the document_chunks
side table, kept in lockstep with every insert into the vec0 table."""
self._conn.executemany(
_INSERT_CHUNK_INDEX,
[(chunk_id, document_id) for chunk_id, document_id, *_ in rows],
)
def _delete_chunks_by_document_id(self, document_id: str) -> None:
"""Delete all of a document's chunks via point-deletes on `id`.
vec0 has no efficient lookup on the document_id metadata column
outside a KNN query (see _open_connection), so a plain
`DELETE ... WHERE document_id = ?` is a full table scan regardless of
index size. Looking the chunk ids up in document_chunks first (a real
indexed lookup) and deleting each by its `id` primary key instead
turns that scan into a handful of O(1) point deletes.
"""
doc_id = str(document_id)
chunk_rows = self._conn.execute(
"SELECT chunk_id FROM document_chunks WHERE document_id = ?",
(doc_id,),
).fetchall()
self._conn.executemany(
"DELETE FROM " + DEFAULT_TABLE_NAME + " WHERE id = ?",
[(row["chunk_id"],) for row in chunk_rows],
)
self._conn.execute(
"DELETE FROM document_chunks WHERE document_id = ?",
(doc_id,),
)
def _increment_total_inserts(self, count: int) -> None:
"""Increment the cumulative insert counter stored in index_meta.
@@ -348,7 +434,8 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
rows = [self._row(node) for node in nodes]
with self._transaction():
self._ensure_table(len(nodes[0].get_embedding()))
self._conn.executemany(self._INSERT, rows)
self._conn.executemany(_INSERT, rows)
self._index_chunks(rows)
self._increment_total_inserts(len(rows))
return [node.node_id for node in nodes]
@@ -365,22 +452,17 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
if nodes:
self._ensure_table(len(nodes[0].get_embedding()))
if self.table_exists():
self._conn.execute(
"DELETE FROM " + DEFAULT_TABLE_NAME + " WHERE document_id = ?",
(str(document_id),),
)
self._delete_chunks_by_document_id(document_id)
if rows:
self._conn.executemany(self._INSERT, rows)
self._conn.executemany(_INSERT, rows)
self._index_chunks(rows)
self._increment_total_inserts(len(rows))
return [node.node_id for node in nodes]
def delete(self, ref_doc_id: str, **delete_kwargs: Any) -> None:
if self.table_exists():
with self._transaction():
self._conn.execute(
"DELETE FROM " + DEFAULT_TABLE_NAME + " WHERE document_id = ?",
(str(ref_doc_id),),
)
self._delete_chunks_by_document_id(ref_doc_id)
def _rows_to_nodes(self, rows: list[sqlite3.Row]) -> list[BaseNode]:
nodes: list[BaseNode] = []
@@ -464,6 +546,67 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
result[doc_id] = str(row["modified"] or "")
return result
@property
def _db_path(self) -> str:
return str(Path(self._uri) / DB_FILENAME)
@contextmanager
def _rebuild_file(self) -> Iterator[sqlite3.Connection]:
"""Open a fresh temp database file for a file-swap rebuild (compact
or structural migration), yielding its connection for the caller to
populate.
On success, swaps the temp file in as the live database (closing
this store's current connection first -- see _swap_in_compact()).
On any exception, discards the temp file, including its -wal/-shm,
instead, and this store's own connection is left untouched.
"""
compact_path = self._db_path + ".compact"
new_conn = self._open_connection(compact_path)
try:
yield new_conn
except BaseException:
new_conn.close()
for suffix in ["", "-wal", "-shm"]:
Path(compact_path + suffix).unlink(missing_ok=True)
raise
else:
new_conn.close()
self._swap_in_compact(compact_path, self._db_path)
@staticmethod
def _rebuild_into(
src_conn: sqlite3.Connection,
dst_conn: sqlite3.Connection,
dim: int,
meta_keys: tuple[str, ...] = ("dim", "embed_model", "schema_version"),
) -> int:
"""Create the vec0 table in ``dst_conn``, copy ``meta_keys`` from
``src_conn``'s index_meta, and stream every live row across
(populating document_chunks as it goes -- see _copy_rows()).
Returns the number of rows copied.
Used by both compact() (default meta_keys: the schema is unchanged)
and structural migrations (meta_keys minus "schema_version", which
the migration sets to its own target version instead of preserving
the source's).
"""
PaperlessSqliteVecVectorStore._create_vec_table(dst_conn, dim)
for key in meta_keys:
value = PaperlessSqliteVecVectorStore._meta_get_on(src_conn, key)
if value is not None:
PaperlessSqliteVecVectorStore._meta_set_on(dst_conn, key, value)
dst_conn.execute("BEGIN IMMEDIATE")
copied = _copy_rows(src_conn, dst_conn)
# Reset the cumulative counter: after a rebuild, total_inserts == live.
PaperlessSqliteVecVectorStore._meta_set_on(
dst_conn,
"total_inserts",
str(copied),
)
dst_conn.execute("COMMIT")
return copied
def compact(self, *, force: bool = False) -> None:
"""Rebuild the database file to reclaim space left behind by DELETEs.
@@ -496,50 +639,8 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
live,
total,
)
db_path = str(Path(self._uri) / DB_FILENAME)
compact_path = db_path + ".compact"
# Copy all live rows into a fresh database file.
new_conn = self._open_connection(compact_path)
try:
self._create_vec_table(new_conn, dim)
self._meta_set_on(new_conn, "dim", str(dim))
for key in ("embed_model", "schema_version"):
value = self._meta_get(key)
if value is not None:
self._meta_set_on(new_conn, key, value)
src_cursor = self._conn.execute(
"SELECT id, document_id, modified, node_content, embedding "
"FROM " + DEFAULT_TABLE_NAME,
)
new_conn.execute("BEGIN IMMEDIATE")
# Stream rows from the source cursor in batches instead of
# materializing the whole table in memory, so a large index does
# not cause an OOM during routine maintenance compactions.
while batch := src_cursor.fetchmany(COMPACT_BATCH_SIZE):
new_conn.executemany(
self._INSERT,
[
(
r["id"],
r["document_id"],
r["modified"],
r["node_content"],
bytes(r["embedding"]),
)
for r in batch
],
)
# Reset the cumulative counter: after compact, total_inserts == live.
self._meta_set_on(new_conn, "total_inserts", str(live))
new_conn.execute("COMMIT")
except BaseException:
new_conn.close()
for p in [compact_path, compact_path + "-wal", compact_path + "-shm"]:
Path(p).unlink(missing_ok=True)
raise
new_conn.close()
self._swap_in_compact(compact_path, db_path)
with self._rebuild_file() as new_conn:
self._rebuild_into(self._conn, new_conn, dim)
def _swap_in_compact(self, compact_path: str, db_path: str) -> None:
"""Atomically replace the live database with the compacted copy."""
@@ -551,6 +652,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 +685,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 +703,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,
@@ -601,16 +725,12 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
dim = self.vector_dim()
if dim is None: # pragma: no cover
raise RuntimeError("Cannot migrate: no stored vector dimension")
db_path = str(Path(self._uri) / DB_FILENAME)
compact_path = db_path + ".compact"
new_conn = self._open_connection(compact_path)
try:
with self._rebuild_file() as new_conn:
migration.apply(self._conn, new_conn, dim)
self._meta_set_on(new_conn, "schema_version", str(migration.to_version))
except BaseException: # pragma: no cover
new_conn.close()
for p in [compact_path, compact_path + "-wal", compact_path + "-shm"]:
Path(p).unlink(missing_ok=True)
raise
new_conn.close()
self._swap_in_compact(compact_path, db_path)
# Registers m0001 into MIGRATIONS; must be at the bottom (needs
# PaperlessSqliteVecVectorStore fully defined) -- see
# paperless_ai/migrations/__init__.py for the full procedure.
from paperless_ai.migrations import m0001_add_document_chunks # noqa: E402, F401