Coverage for paperless_ai/vector_store.py: 0%
328 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
1import json
2import logging
3import sqlite3
4import struct
5from collections.abc import Iterator
6from collections.abc import Sequence
7from contextlib import contextmanager
8from pathlib import Path
9from types import TracebackType
10from typing import Any
11from typing import NamedTuple
13import sqlite_vec
14from llama_index.core.bridge.pydantic import PrivateAttr
15from llama_index.core.schema import BaseNode
16from llama_index.core.vector_stores.types import BasePydanticVectorStore
17from llama_index.core.vector_stores.types import FilterCondition
18from llama_index.core.vector_stores.types import FilterOperator
19from llama_index.core.vector_stores.types import MetadataFilter
20from llama_index.core.vector_stores.types import MetadataFilters
21from llama_index.core.vector_stores.types import VectorStoreQuery
22from llama_index.core.vector_stores.types import VectorStoreQueryResult
23from llama_index.core.vector_stores.utils import metadata_dict_to_node
24from llama_index.core.vector_stores.utils import node_to_metadata_dict
26from paperless_ai.migrations import MIGRATIONS
27from paperless_ai.migrations import Migration
28from paperless_ai.tables import ChunkRow
29from paperless_ai.tables import DocumentChunksTable
30from paperless_ai.tables import DocumentMetaRow
31from paperless_ai.tables import DocumentMetaTable
32from paperless_ai.tables import IndexMetaTable
34logger = logging.getLogger("paperless_ai.vector_store")
36DB_FILENAME = "llmindex.db"
37DEFAULT_TABLE_NAME = "documents"
39# Current schema version. Written to index_meta at table creation and bumped
40# whenever a Migration is added to MIGRATIONS. check_and_run_migrations() uses
41# this to decide which migrations to run on an existing store.
42SCHEMA_VERSION = 2
44# compact(): rebuild when the cumulative rowid count exceeds this multiple of
45# the live row count. DELETEs on vec0 tables never reclaim space (upstream
46# asg017/sqlite-vec#54), so per-document re-index churn grows the file until
47# a rebuild copies the live rows into a fresh table.
48COMPACT_BLOAT_RATIO = 2.0
50# Number of rows fetched/copied per batch whenever this module streams rows
51# instead of materializing them all at once, keeping memory bounded regardless
52# of index size -- used by compact()'s rebuild, m0001_v1_to_v2's migration
53# copy, and DocumentMetaTable.copy_all(). No longer compact()-specific, hence
54# the plain name.
55BATCH_SIZE = 500
57# Filterable vec0 metadata columns. _build_where() only ever receives filter
58# keys we construct ourselves, but allowlisting keeps SQL identifiers safe by
59# construction. "modified" is not here: it is never filtered on, and as of
60# schema v2 it isn't even a vec0 column anymore (see document_meta).
61_FILTER_COLUMNS = frozenset({"document_id"})
64class _Row(NamedTuple):
65 """One node, ready to write. ``modified`` is not a vec0 column (see
66 document_meta) -- it rides along here because every row-producing call
67 site needs both the vec0 insert values and the document_meta upsert
68 value from the same node.
69 """
71 chunk_id: str
72 document_id: int
73 modified: str
74 node_content: str
75 embedding: bytes
78# _build_where(): the largest IN value list translated into bound SQL
79# parameters. SQLite's own hard limit (SQLITE_MAX_VARIABLE_NUMBER) is 32766
80# by default; this leaves headroom below that for the query's other bound
81# parameters (the embedding blob, k, and any NE clause) and for the limit
82# itself to move. An IN filter this large should not happen in practice --
83# callers are expected to pass None (no filter) rather than every id when
84# the filter would not actually narrow anything -- so this is a guard
85# against a future regression, not a normal code path.
86_MAX_IN_VALUES = 32700
89def _pack(embedding: Sequence[float]) -> bytes:
90 return struct.pack(f"{len(embedding)}f", *embedding)
93def _unpack(blob: bytes) -> list[float]:
94 return list(struct.unpack(f"{len(blob) // 4}f", blob))
97_INSERT = (
98 "INSERT INTO "
99 + DEFAULT_TABLE_NAME
100 + " (id, document_id, node_content, embedding) VALUES (?, ?, ?, ?)"
101)
104def _vec0_params(rows: list[_Row]) -> list[tuple[str, int, str, bytes]]:
105 """``rows``, minus the ``modified`` field vec0 no longer stores."""
106 return [(r.chunk_id, r.document_id, r.node_content, r.embedding) for r in rows]
109def _build_where(filters: MetadataFilters | None) -> tuple[str, list[int]]:
110 """Translate the EQ / IN / NIN / NE filters we use into a parameterized
111 SQL clause on vec0 metadata columns. Returns ("", []) when there is
112 nothing to filter. document_id is vec0's only filterable column and is
113 INTEGER; every value is coerced via int() here so callers (which today
114 still pass strings in places, e.g. indexing.py's MetadataFilter
115 construction) don't have to be individually correct -- vec0 doesn't
116 coerce types itself.
117 """
118 if filters is None or not filters.filters:
119 return "", []
120 clauses: list[str] = []
121 params: list[int] = []
122 for f in filters.filters:
123 # filters.filters is Union[MetadataFilter, ExactMatchFilter, MetadataFilters];
124 # we only build MetadataFilter entries, so skip anything else at runtime.
125 if not isinstance(f, MetadataFilter):
126 continue
127 if f.key not in _FILTER_COLUMNS: # pragma: no cover - we build the keys
128 raise NotImplementedError(f"Unsupported filter column: {f.key}")
129 if f.operator in (FilterOperator.IN, FilterOperator.NIN):
130 is_in = f.operator == FilterOperator.IN
131 sql_op = "IN" if is_in else "NOT IN"
132 values = [int(v) for v in f.value] # type: ignore[union-attr]
133 if not values:
134 # An empty IN list matches nothing; an empty NOT IN list
135 # excludes nothing, so it matches everything.
136 clauses.append("1 = 0" if is_in else "1 = 1")
137 continue
138 if len(values) > _MAX_IN_VALUES:
139 # Refuse rather than risk SQLite's own bound-parameter limit
140 # ("too many SQL variables"): a list this large must match no
141 # rows, never widen the scope to "everything" -- true for
142 # NOT IN too, where failing open would surface every
143 # excluded row.
144 logger.warning(
145 "Refusing to build a %s filter on %r with %d values "
146 "(over the %d-value safety limit); returning no rows.",
147 sql_op,
148 f.key,
149 len(values),
150 _MAX_IN_VALUES,
151 )
152 clauses.append("1 = 0")
153 continue
154 placeholders = ",".join("?" for _ in values)
155 clauses.append(f"{f.key} {sql_op} ({placeholders})")
156 params.extend(values)
157 elif f.operator == FilterOperator.EQ:
158 clauses.append(f"{f.key} = ?")
159 params.append(int(f.value))
160 elif f.operator == FilterOperator.NE:
161 clauses.append(f"{f.key} != ?")
162 params.append(int(f.value))
163 else: # pragma: no cover - we only ever build EQ/IN/NIN/NE filters
164 raise NotImplementedError(f"Unsupported filter operator: {f.operator}")
165 if not clauses:
166 # Filters were requested but none could be translated. Fail closed
167 # rather than emit "()" (invalid SQL): filters scope document access,
168 # so an empty translation must match no rows, never widen the scope.
169 return "1 = 0", []
170 joiner = " OR " if filters.condition == FilterCondition.OR else " AND "
171 return "(" + joiner.join(clauses) + ")", params
174class PaperlessSqliteVecVectorStore(BasePydanticVectorStore):
175 """A llama-index vector store backed by a sqlite-vec vec0 table.
177 Stores one row per node: the node id (TEXT primary key), its document id
178 (metadata column, used for EQ/IN filtering and per-document delete), the
179 document's modified timestamp, the embedding (float32, cosine metric), and
180 the serialized node (text + metadata) as JSON in an auxiliary column.
181 ``stores_text`` lets llama-index run off this store alone, with no
182 separate docstore or index store.
184 Everything lives in one SQLite database file (``DB_FILENAME``) inside the
185 directory given as ``uri`` (kept as a directory for compatibility with the
186 previous LanceDB layout). WAL mode allows readers in other processes to
187 proceed while the (FileLock-serialized) writer holds a transaction.
189 Implemented surface of ``BasePydanticVectorStore``
190 ---------------------------------------------------
191 Only the methods actively used by this codebase are implemented.
192 ``delete_nodes`` and the ``node_ids`` lookup path of ``get_nodes`` are
193 part of the llama-index interface contract and may be needed if a future
194 retriever or extension invokes them — add them then, with tests.
195 """
197 stores_text: bool = True
198 flat_metadata: bool = False
200 _uri: str = PrivateAttr()
201 _embed_model_name: str | None = PrivateAttr()
202 _conn: Any = PrivateAttr()
204 def __init__(
205 self,
206 uri: str,
207 embed_model_name: str | None = None,
208 ) -> None:
209 super().__init__(stores_text=True, flat_metadata=False)
210 self._uri = uri
211 self._embed_model_name = embed_model_name
212 self._conn = self._open_connection(str(Path(uri) / DB_FILENAME))
214 @staticmethod
215 def _open_connection(db_path: str) -> sqlite3.Connection:
216 conn = sqlite3.connect(
217 db_path,
218 timeout=30,
219 isolation_level=None, # autocommit; explicit transactions below
220 )
221 conn.row_factory = sqlite3.Row
222 conn.enable_load_extension(True) # noqa: FBT003
223 sqlite_vec.load(conn)
224 conn.enable_load_extension(False) # noqa: FBT003
225 conn.execute("PRAGMA journal_mode=WAL")
226 conn.execute("PRAGMA synchronous=NORMAL")
227 IndexMetaTable.create(conn)
228 # vec0 metadata columns only get an efficient lookup path inside a
229 # KNN (MATCH) query; a plain `WHERE document_id = ?` is a full table
230 # scan regardless of index size. This plain, indexed table is how
231 # delete()/upsert_document() find a document's chunk ids without
232 # that scan.
233 DocumentChunksTable.create(conn)
234 # modified used to be a vec0 metadata column, but vec0 only inlines
235 # TEXT metadata up to 12 bytes -- an ISO timestamp is always longer,
236 # so every read recompiled and stepped a fresh SQL statement per row.
237 # It was never filtered on inside a KNN query either, so it never
238 # needed to be a vec0 column at all. One row per document here (not
239 # per chunk, like document_chunks), since every chunk of a document
240 # shares the same modified value -- see get_modified_times().
241 DocumentMetaTable.create(conn)
242 return conn
244 @property
245 def client(self) -> Any:
246 return self._conn
248 def close(self) -> None:
249 """Close the underlying SQLite connection (idempotent)."""
250 self._conn.close()
252 def __enter__(self) -> "PaperlessSqliteVecVectorStore":
253 return self
255 def __exit__(
256 self,
257 exc_type: type[BaseException] | None,
258 exc_val: BaseException | None,
259 exc_tb: TracebackType | None,
260 ) -> None:
261 # Deterministically release the connection (and its WAL/SHM handles) so
262 # it is never left open across a compaction/migration file swap.
263 self.close()
265 @contextmanager
266 def _transaction(self) -> Iterator[None]:
267 self._conn.execute("BEGIN IMMEDIATE")
268 try:
269 yield
270 except BaseException: # pragma: no cover
271 self._conn.execute("ROLLBACK")
272 raise
273 else:
274 self._conn.execute("COMMIT")
276 def table_exists(self) -> bool:
277 return (
278 self._conn.execute(
279 "SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?",
280 (DEFAULT_TABLE_NAME,),
281 ).fetchone()
282 is not None
283 )
285 def vector_dim(self) -> int | None:
286 if not self.table_exists():
287 return None
288 return IndexMetaTable.get_dim(self._conn)
290 def drop_table(self) -> None:
291 self._conn.execute("DROP TABLE IF EXISTS " + DEFAULT_TABLE_NAME)
292 self._conn.execute("DELETE FROM index_meta")
293 DocumentChunksTable.delete_all(self._conn)
294 DocumentMetaTable.delete_all(self._conn)
296 def stored_model_name(self) -> str | None:
297 """Return the embedding model name recorded at table creation, or None."""
298 if not self.table_exists():
299 return None
300 return IndexMetaTable.get_embed_model(self._conn)
302 def config_mismatch(self, model_name: str) -> bool:
303 """True when the stored model name differs from ``model_name``.
305 Returns False when no table exists or when the table predates
306 model-name tracking — conservative default avoids spurious rebuilds.
307 """
308 stored = self.stored_model_name()
309 if stored is None:
310 return False
311 return stored != model_name
313 @staticmethod
314 def _create_vec_table(conn: sqlite3.Connection, dim: int) -> None:
315 # document_id is deliberately a metadata column, NOT a partition key:
316 # partition keys change KNN `k` to per-partition semantics under IN
317 # filters (asg017/sqlite-vec#142); metadata columns give a correct
318 # global top-k. INTEGER (not TEXT, as in schema v1): EQ/NE/IN
319 # comparisons become a native i64 array compare instead of per-row
320 # strncmp against a 16-byte text view, and this drops the unused
321 # metadatatext shadow table TEXT columns carry. modified is not a
322 # column here at all as of v2 -- see document_meta.
323 conn.execute( # nosemgrep: python.sqlalchemy.security.sqlalchemy-execute-raw-query.sqlalchemy-execute-raw-query
324 "CREATE VIRTUAL TABLE "
325 + DEFAULT_TABLE_NAME
326 + " USING vec0("
327 + "id TEXT PRIMARY KEY,"
328 + " document_id INTEGER,"
329 + " +node_content TEXT,"
330 + " embedding float["
331 + str(int(dim))
332 + "] distance_metric=cosine"
333 + ")",
334 )
336 def _create_table(self, dim: int) -> None:
337 self._create_vec_table(self._conn, dim)
338 IndexMetaTable.set_dim(self._conn, dim)
339 IndexMetaTable.set_schema_version(self._conn, SCHEMA_VERSION)
340 if self._embed_model_name:
341 IndexMetaTable.set_embed_model(self._conn, self._embed_model_name)
343 def _ensure_table(self, dim: int, *, table_exists: bool) -> None:
344 if not table_exists:
345 self._create_table(dim)
347 def _row(self, node: BaseNode) -> _Row:
348 meta = node_to_metadata_dict(
349 node,
350 remove_text=False,
351 flat_metadata=self.flat_metadata,
352 )
353 document_id = node.ref_doc_id or node.metadata.get("document_id")
354 return _Row(
355 chunk_id=node.node_id,
356 # document_id is required -- int(None) raises TypeError and
357 # int("not-a-number") raises ValueError, both intentional:
358 # fail loudly on a malformed/missing document_id rather than
359 # silently indexing a chunk with no owning document. modified,
360 # below, still uses the str(x or "") sentinel pattern because a
361 # missing modified value is legitimate (vec0 no longer even
362 # stores it -- see document_meta), whereas document_id must
363 # always be present.
364 document_id=int(document_id),
365 modified=str(node.metadata.get("modified") or ""),
366 node_content=json.dumps(meta),
367 embedding=_pack(node.get_embedding()),
368 )
370 def _index_chunks(self, rows: list[_Row]) -> None:
371 """Record each row's (chunk_id, document_id) in document_chunks, and
372 each row's (document_id, modified) in document_meta -- deduped
373 within the batch, since every chunk of a document shares the same
374 modified value -- kept in lockstep with every insert into the vec0
375 table.
376 """
377 DocumentChunksTable.insert_many(
378 self._conn,
379 (ChunkRow(r.chunk_id, r.document_id) for r in rows),
380 )
381 modified_by_document = {r.document_id: r.modified for r in rows}
382 DocumentMetaTable.upsert_many(
383 self._conn,
384 (
385 DocumentMetaRow(doc_id, mod)
386 for doc_id, mod in modified_by_document.items()
387 ),
388 )
390 def _delete_chunks_by_document_id(self, document_id: int) -> None:
391 """Delete all of a document's chunks via point-deletes on `id`.
393 vec0 has no efficient lookup on the document_id metadata column
394 outside a KNN query, so a plain `DELETE ... WHERE document_id = ?`
395 is a full table scan regardless of index size. Looking the chunk
396 ids up in document_chunks first (a real indexed lookup) and
397 deleting each by its `id` primary key instead turns that scan into
398 a handful of O(1) point deletes.
399 """
400 chunk_ids = DocumentChunksTable.chunk_ids_for_document(
401 self._conn,
402 document_id,
403 )
404 self._conn.executemany(
405 "DELETE FROM " + DEFAULT_TABLE_NAME + " WHERE id = ?",
406 [(chunk_id,) for chunk_id in chunk_ids],
407 )
408 DocumentChunksTable.delete_for_document(self._conn, document_id)
409 DocumentMetaTable.delete_for_document(self._conn, document_id)
411 def _increment_total_inserts(self, count: int) -> None:
412 """Increment the cumulative insert counter stored in index_meta.
414 This counter never decreases (DELETEs do not decrement it) and is
415 used by compact() to estimate the bloat ratio: when total_inserts /
416 live_rows exceeds COMPACT_BLOAT_RATIO the table has accumulated
417 enough deleted-but-not-freed rows to warrant a rebuild.
418 """
419 IndexMetaTable.increment_total_inserts(self._conn, count)
421 def add(self, nodes: Sequence[BaseNode], **add_kwargs: Any) -> list[str]:
422 if not nodes:
423 return []
424 rows = [self._row(node) for node in nodes]
425 with self._transaction():
426 self._ensure_table(
427 len(nodes[0].get_embedding()),
428 table_exists=self.table_exists(),
429 )
430 self._conn.executemany(_INSERT, _vec0_params(rows))
431 self._index_chunks(rows)
432 self._increment_total_inserts(len(rows))
433 return [node.node_id for node in nodes]
435 def upsert_document(
436 self,
437 document_id: int | str,
438 nodes: list[BaseNode],
439 ) -> list[str]:
440 """Atomically replace all stored chunks of ``document_id`` with ``nodes``.
442 One transaction deletes the document's existing rows and inserts the
443 new set (vec0's INSERT OR REPLACE is broken upstream, so delete+insert
444 it is). WAL readers in other processes see either the old or the new
445 chunk set, never a partial state.
446 """
447 doc_id = int(document_id)
448 rows = [self._row(node) for node in nodes]
449 with self._transaction():
450 table_exists = self.table_exists()
451 if nodes and not table_exists:
452 self._ensure_table(
453 len(nodes[0].get_embedding()),
454 table_exists=False,
455 )
456 table_exists = True
457 if table_exists:
458 self._delete_chunks_by_document_id(doc_id)
459 if rows:
460 self._conn.executemany(_INSERT, _vec0_params(rows))
461 self._index_chunks(rows)
462 self._increment_total_inserts(len(rows))
463 return [node.node_id for node in nodes]
465 def delete(self, ref_doc_id: int | str, **delete_kwargs: Any) -> None:
466 if self.table_exists():
467 with self._transaction():
468 self._delete_chunks_by_document_id(int(ref_doc_id))
470 def _rows_to_nodes(self, rows: list[sqlite3.Row]) -> list[BaseNode]:
471 nodes: list[BaseNode] = []
472 for row in rows:
473 node = metadata_dict_to_node(json.loads(row["node_content"]))
474 node.embedding = _unpack(row["embedding"])
475 nodes.append(node)
476 return nodes
478 def get_nodes(
479 self,
480 node_ids: list[str] | None = None,
481 filters: MetadataFilters | None = None,
482 **kwargs: Any,
483 ) -> list[BaseNode]:
484 if node_ids is not None: # pragma: no cover
485 # node_ids lookup is not implemented; see class docstring.
486 raise NotImplementedError(
487 "PaperlessSqliteVecVectorStore does not support node_ids lookup",
488 )
489 if not self.table_exists():
490 return []
491 where, params = _build_where(filters)
492 sql = "SELECT node_content, embedding FROM " + DEFAULT_TABLE_NAME
493 if where:
494 sql += " WHERE " + where
495 return self._rows_to_nodes(self._conn.execute(sql, params).fetchall())
497 def query(
498 self,
499 query: VectorStoreQuery,
500 **kwargs: Any,
501 ) -> VectorStoreQueryResult:
502 if not self.table_exists():
503 return VectorStoreQueryResult(nodes=[], similarities=[], ids=[])
504 if query.query_embedding is None: # pragma: no cover
505 return VectorStoreQueryResult(nodes=[], similarities=[], ids=[])
506 top_k = query.similarity_top_k if query.similarity_top_k is not None else 10
507 where, params = _build_where(query.filters)
508 sql = (
509 "SELECT id, node_content, embedding, distance FROM "
510 + DEFAULT_TABLE_NAME
511 + " WHERE embedding MATCH ? AND k = ?"
512 )
513 if where:
514 sql += " AND " + where
515 rows = self._conn.execute(
516 sql,
517 [_pack(query.query_embedding), top_k, *params],
518 ).fetchall()
519 # vec0 returns rows distance-sorted ascending; slice defensively in
520 # case future schema changes alter k semantics (e.g. partition keys
521 # return k rows per partition).
522 rows = rows[:top_k]
523 nodes = self._rows_to_nodes(rows)
524 # Cosine distance in [0, 2]; map to a descending similarity.
525 # vec0 returns None distance when the query embedding is the zero vector
526 # (no meaningful cosine angle); treat that as maximum distance (1.0) so
527 # the row is included but ranked last.
528 sims = [
529 1.0 - float(row["distance"] if row["distance"] is not None else 1.0)
530 for row in rows
531 ]
532 ids = [row["id"] for row in rows]
533 return VectorStoreQueryResult(nodes=nodes, similarities=sims, ids=ids)
535 def get_modified_times(self) -> dict[str, str]:
536 """Return {document_id: stored_modified_isoformat} for all indexed documents.
538 document_meta already has exactly one row per document (not per
539 chunk, unlike the vec0 table), so no dedup is needed here.
540 """
541 if not self.table_exists():
542 return {}
543 return DocumentMetaTable.all_modified_times(self._conn)
545 @property
546 def _db_path(self) -> str:
547 return str(Path(self._uri) / DB_FILENAME)
549 @contextmanager
550 def _rebuild_file(self) -> Iterator[sqlite3.Connection]:
551 """Open a fresh temp database file for a file-swap rebuild (compact
552 or structural migration), yielding its connection for the caller to
553 populate.
555 On success, swaps the temp file in as the live database (closing
556 this store's current connection first -- see _swap_in_compact()).
557 On any exception, discards the temp file, including its -wal/-shm,
558 instead, and this store's own connection is left untouched.
559 """
560 compact_path = self._db_path + ".compact"
561 new_conn = self._open_connection(compact_path)
562 try:
563 yield new_conn
564 except BaseException:
565 new_conn.close()
566 for suffix in ["", "-wal", "-shm"]:
567 Path(compact_path + suffix).unlink(missing_ok=True)
568 raise
569 else:
570 new_conn.close()
571 self._swap_in_compact(compact_path, self._db_path)
573 def compact(self, *, force: bool = False) -> None:
574 """Rebuild the database file to reclaim space left behind by DELETEs.
576 vec0 DELETE only invalidates rows; the vector data stays in the file
577 forever, and per-document re-indexing is a delete+insert. The
578 cumulative insert counter in ``index_meta`` tracks total rows ever
579 written; when that exceeds ``COMPACT_BLOAT_RATIO`` x the live row
580 count (or when forced), live rows are copied into a fresh database
581 file and swapped in via ``os.replace``.
583 Note: ``ALTER TABLE ... RENAME TO`` on vec0 virtual tables does NOT
584 rename the shadow tables (sqlite-vec upstream limitation), so an
585 in-place rename-based rebuild is not safe. The file-swap approach is
586 the maintainer-endorsed workaround.
587 """
588 if not self.table_exists():
589 return
590 if self.has_pending_migration():
591 logger.warning(
592 "Skipping compact: store has a pending schema migration; "
593 "run check_and_run_migrations() first",
594 )
595 return
596 live = DocumentChunksTable.count(self._conn)
597 total = IndexMetaTable.get_total_inserts(self._conn) or live
598 if not force and total <= max(live, 1) * COMPACT_BLOAT_RATIO:
599 return
600 dim = self.vector_dim()
601 if dim is None: # pragma: no cover - dim is written at creation
602 logger.warning("Skipping compact: no stored vector dimension")
603 return
604 logger.info(
605 "Compacting LLM index (%d live rows, %d cumulative inserts)",
606 live,
607 total,
608 )
609 with self._rebuild_file() as new_conn:
610 self._rebuild_into(self._conn, new_conn, dim)
612 @staticmethod
613 def _rebuild_into(
614 src_conn: sqlite3.Connection,
615 dst_conn: sqlite3.Connection,
616 dim: int,
617 ) -> None:
618 """Create the vec0 table in ``dst_conn``, copy dim/embed_model from
619 ``src_conn``, and stream every live vec0 row, document_chunks row,
620 and document_meta row across. Used by compact() only --
621 m0001_v1_to_v2 freezes its own copy loop instead of calling this,
622 since this always reflects the *current* schema (see the migration
623 DDL-freezing rule in the spec).
624 """
625 PaperlessSqliteVecVectorStore._create_vec_table(dst_conn, dim)
626 dim_value = IndexMetaTable.get_dim(src_conn)
627 if dim_value is not None:
628 IndexMetaTable.set_dim(dst_conn, dim_value)
629 embed_model = IndexMetaTable.get_embed_model(src_conn)
630 if embed_model is not None:
631 IndexMetaTable.set_embed_model(dst_conn, embed_model)
632 schema_version = IndexMetaTable.get_schema_version(src_conn)
633 if schema_version is not None:
634 IndexMetaTable.set_schema_version(dst_conn, schema_version)
636 dst_conn.execute("BEGIN IMMEDIATE")
637 src_cursor = src_conn.execute(
638 "SELECT id, document_id, node_content, embedding FROM "
639 + DEFAULT_TABLE_NAME,
640 )
641 copied = 0
642 while batch := src_cursor.fetchmany(BATCH_SIZE):
643 dst_conn.executemany(
644 _INSERT,
645 [
646 (
647 r["id"],
648 r["document_id"],
649 r["node_content"],
650 bytes(r["embedding"]),
651 )
652 for r in batch
653 ],
654 )
655 DocumentChunksTable.insert_many(
656 dst_conn,
657 (ChunkRow(r["id"], r["document_id"]) for r in batch),
658 )
659 copied += len(batch)
660 DocumentMetaTable.copy_all(src_conn, dst_conn, BATCH_SIZE)
661 # Reset the cumulative counter: after a rebuild, total_inserts == live.
662 IndexMetaTable.reset_total_inserts(dst_conn, copied)
663 dst_conn.execute("COMMIT")
665 def _swap_in_compact(self, compact_path: str, db_path: str) -> None:
666 """Atomically replace the live database with the compacted copy."""
667 self._conn.close()
668 for suffix in ["-wal", "-shm"]:
669 stale = Path(compact_path + suffix)
670 if stale.exists(): # pragma: no cover
671 stale.unlink()
672 Path(compact_path).replace(db_path)
673 self._conn = self._open_connection(db_path)
675 def _stored_schema_version(self) -> int | None:
676 """The schema_version recorded in index_meta, or None if no table
677 exists. A missing key (a store predating version tracking) is
678 treated as SCHEMA_VERSION -- i.e. already current -- since no
679 migration in MIGRATIONS targets a version before tracking began.
680 """
681 if not self.table_exists():
682 return None
683 raw_version = IndexMetaTable.get_schema_version(self._conn)
684 return raw_version if raw_version is not None else SCHEMA_VERSION
686 def has_pending_migration(self) -> bool:
687 """Cheaply check whether a migration is pending, with no exclusive
688 access needed -- just a metadata read under the connection callers
689 already hold via the write FileLock.
691 Callers should only pay for check_and_run_migrations()'s exclusive
692 access (a structural migration's file swap must not run while
693 readers are active) when this returns True, so that the common
694 case -- already at SCHEMA_VERSION -- never contends with readers
695 or a concurrent compaction.
696 """
697 current = self._stored_schema_version()
698 return current is not None and current < SCHEMA_VERSION
700 def check_and_run_migrations(self) -> bool:
701 """Apply any pending schema migrations to the store.
703 Structural migrations copy live rows into a new-schema file with no
704 re-embedding. Re-embed migrations cannot be applied automatically;
705 this method returns True when one is encountered so the caller can
706 force a full rebuild (which recreates the table at SCHEMA_VERSION).
708 Must be called under the write FileLock, with readers excluded (see
709 has_pending_migration() for a cheap pre-check that avoids paying for
710 that exclusion in the common case). No-op when the table does not
711 exist or is already at SCHEMA_VERSION.
712 """
713 current = self._stored_schema_version()
714 if current is None or current >= SCHEMA_VERSION:
715 return False
717 pending = sorted(
718 [m for m in MIGRATIONS if current <= m.from_version < SCHEMA_VERSION],
719 key=lambda m: m.from_version,
720 )
722 for migration in pending:
723 if migration.kind == "re-embed":
724 logger.warning(
725 "LLM index schema v%d -> v%d requires re-embedding (%s); "
726 "the caller must force a rebuild.",
727 migration.from_version,
728 migration.to_version,
729 migration.description,
730 )
731 return True
732 logger.info(
733 "Running structural LLM index migration v%d -> v%d: %s",
734 migration.from_version,
735 migration.to_version,
736 migration.description,
737 )
738 self._run_structural_migration(migration)
740 return False
742 def _run_structural_migration(self, migration: Migration) -> None:
743 """Execute a structural migration using the same file-swap as compact()."""
744 assert migration.apply is not None, "structural migration must have apply()"
745 dim = self.vector_dim()
746 if dim is None: # pragma: no cover
747 raise RuntimeError("Cannot migrate: no stored vector dimension")
748 with self._rebuild_file() as new_conn:
749 migration.apply(self._conn, new_conn, dim)
750 IndexMetaTable.set_schema_version(new_conn, migration.to_version)
753# Registers m0001_v1_to_v2 into MIGRATIONS; must be at the bottom (needs
754# PaperlessSqliteVecVectorStore fully defined) -- see
755# paperless_ai/migrations/__init__.py for the full procedure.
756from paperless_ai.migrations import m0001_v1_to_v2 # noqa: E402, F401