Compare commits

...
Author SHA1 Message Date
stumpylogandClaude Sonnet 5 feac526318 Perf: count compact()'s live rows via document_chunks, not a vec0 scan
compact() ran SELECT count(*) on the vec0 table itself purely to compute
the bloat ratio, on every update_llm_index() call -- including small
scoped ones from bulk_update_documents. Like any other document_id-scale
lookup, count(*) has no point-lookup fast path on vec0's own table, so
this is a full scan regardless of index size.

document_chunks is kept in lockstep with the vec0 table by construction
and is a plain indexed table, so it gives an identical count far more
cheaply. Only affects the bloat-ratio heuristic and a log line -- the
actual rebuild already streams and counts real vec0 rows independently
(_copy_rows()'s own return value), so this can't desync compaction's
correctness even in the brief window before a store's document_chunks is
fully backfilled by migration.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 15:36:42 -07:00
stumpylogandClaude Sonnet 5 3838706194 Fix: freeze m0001's vec0 table shape instead of tracking current schema
_rebuild_into()/_create_vec_table() always reflect whatever the *current*
schema is -- correct for compact() (which only ever runs on an
already-current-schema store), but wrong for a migration that isn't the
latest one anymore. m0001 delegated to them, which worked fine as long as
v2 was current, but silently breaks the moment a v2 -> v3 migration is
added: m0001 would then produce a v3-shaped table (whatever columns
happen to be current) instead of its actual v2 shape, and the next
migration in the chain would find the columns it expects to migrate from
already gone.

Spell out m0001's own v2-shaped CREATE VIRTUAL TABLE and row copy instead,
so it keeps producing the same output forever, independent of later
schema changes -- the same reasoning already documented on the test
helper _copying_apply(), just not, until now, applied to the real
migration.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-28 13:34:25 -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
14 changed files with 1627 additions and 371 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,
+78 -28
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.
@@ -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,95 @@
import sqlite3
from paperless_ai.migrations import MIGRATIONS
from paperless_ai.migrations import Migration
from paperless_ai.vector_store import COMPACT_BATCH_SIZE
from paperless_ai.vector_store import DEFAULT_TABLE_NAME
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.
Deliberately spells out its own v2-shaped vec0 table and row copy,
rather than delegating to PaperlessSqliteVecVectorStore's
_create_vec_table()/_rebuild_into()/_copy_rows(): those always reflect
whatever the *current* schema is. If a later migration changes that
schema (bumping SCHEMA_VERSION again), this migration must keep
producing its own historical v2 shape regardless -- otherwise a user
upgrading across multiple versions in one go (e.g. v1 straight to v4)
would have this migration silently produce a v4-shaped table instead
of v2, and the v2 -> v3 migration that runs right after it would find
the columns it expects to migrate *from* already gone.
"""
dst_conn.execute( # nosemgrep: python.sqlalchemy.security.sqlalchemy-execute-raw-query.sqlalchemy-execute-raw-query
"CREATE VIRTUAL TABLE "
+ DEFAULT_TABLE_NAME
+ " USING vec0("
+ "id TEXT PRIMARY KEY,"
+ " document_id TEXT,"
+ " modified TEXT,"
+ " +node_content TEXT,"
+ " embedding float["
+ str(int(dim))
+ "] distance_metric=cosine"
+ ")",
)
PaperlessSqliteVecVectorStore._meta_set_on(dst_conn, "dim", str(dim))
embed_model = PaperlessSqliteVecVectorStore._meta_get_on(src_conn, "embed_model")
if embed_model is not None:
PaperlessSqliteVecVectorStore._meta_set_on(dst_conn, "embed_model", embed_model)
dst_conn.execute("BEGIN IMMEDIATE")
src_cursor = src_conn.execute(
"SELECT id, document_id, modified, node_content, embedding FROM "
+ DEFAULT_TABLE_NAME,
)
live = 0
while batch := src_cursor.fetchmany(COMPACT_BATCH_SIZE):
dst_conn.executemany(
"INSERT INTO "
+ DEFAULT_TABLE_NAME
+ " (id, document_id, modified, node_content, embedding) "
"VALUES (?, ?, ?, ?, ?)",
[
(
r["id"],
r["document_id"],
r["modified"],
r["node_content"],
bytes(r["embedding"]),
)
for r in batch
],
)
dst_conn.executemany(
"INSERT INTO document_chunks (chunk_id, document_id) VALUES (?, ?)",
[(r["id"], r["document_id"]) for r in batch],
)
live += len(batch)
# This migration only ever copies live rows (like compact()), so the
# cumulative counter resets to match -- the new file has no bloat yet.
PaperlessSqliteVecVectorStore._meta_set_on(dst_conn, "total_inserts", str(live))
dst_conn.execute("COMMIT")
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,
),
)
+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
+222 -124
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
@@ -53,38 +61,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)
@@ -93,6 +69,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.
@@ -189,6 +201,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 +253,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 +293,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 +360,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 +409,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 +427,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 +521,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.
@@ -481,9 +599,12 @@ class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
"""
if not self.table_exists():
return
live = self._conn.execute(
"SELECT count(*) FROM " + DEFAULT_TABLE_NAME,
).fetchone()[0]
# document_chunks is a plain indexed table kept in lockstep with the
# vec0 table (see _index_chunks()/_delete_chunks_by_document_id()),
# so it gives an identical live-row count to `SELECT count(*) FROM
# documents` without vec0's full table scan -- count(*) has no
# equivalent point-lookup fast path, unlike an EQ constraint on `id`.
live = self._conn.execute("SELECT count(*) FROM document_chunks").fetchone()[0]
total = int(self._meta_get("total_inserts") or str(live))
if not force and total <= max(live, 1) * COMPACT_BLOAT_RATIO:
return
@@ -496,50 +617,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 +630,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 +663,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 +681,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 +703,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