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

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 

12 

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 

25 

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 

33 

34logger = logging.getLogger("paperless_ai.vector_store") 

35 

36DB_FILENAME = "llmindex.db" 

37DEFAULT_TABLE_NAME = "documents" 

38 

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 

43 

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 

49 

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 

56 

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"}) 

62 

63 

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 """ 

70 

71 chunk_id: str 

72 document_id: int 

73 modified: str 

74 node_content: str 

75 embedding: bytes 

76 

77 

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 

87 

88 

89def _pack(embedding: Sequence[float]) -> bytes: 

90 return struct.pack(f"{len(embedding)}f", *embedding) 

91 

92 

93def _unpack(blob: bytes) -> list[float]: 

94 return list(struct.unpack(f"{len(blob) // 4}f", blob)) 

95 

96 

97_INSERT = ( 

98 "INSERT INTO " 

99 + DEFAULT_TABLE_NAME 

100 + " (id, document_id, node_content, embedding) VALUES (?, ?, ?, ?)" 

101) 

102 

103 

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] 

107 

108 

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 

172 

173 

174class PaperlessSqliteVecVectorStore(BasePydanticVectorStore): 

175 """A llama-index vector store backed by a sqlite-vec vec0 table. 

176 

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. 

183 

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. 

188 

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 """ 

196 

197 stores_text: bool = True 

198 flat_metadata: bool = False 

199 

200 _uri: str = PrivateAttr() 

201 _embed_model_name: str | None = PrivateAttr() 

202 _conn: Any = PrivateAttr() 

203 

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)) 

213 

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 

243 

244 @property 

245 def client(self) -> Any: 

246 return self._conn 

247 

248 def close(self) -> None: 

249 """Close the underlying SQLite connection (idempotent).""" 

250 self._conn.close() 

251 

252 def __enter__(self) -> "PaperlessSqliteVecVectorStore": 

253 return self 

254 

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() 

264 

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") 

275 

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 ) 

284 

285 def vector_dim(self) -> int | None: 

286 if not self.table_exists(): 

287 return None 

288 return IndexMetaTable.get_dim(self._conn) 

289 

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) 

295 

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) 

301 

302 def config_mismatch(self, model_name: str) -> bool: 

303 """True when the stored model name differs from ``model_name``. 

304 

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 

312 

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 ) 

335 

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) 

342 

343 def _ensure_table(self, dim: int, *, table_exists: bool) -> None: 

344 if not table_exists: 

345 self._create_table(dim) 

346 

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 ) 

369 

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 ) 

389 

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`. 

392 

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) 

410 

411 def _increment_total_inserts(self, count: int) -> None: 

412 """Increment the cumulative insert counter stored in index_meta. 

413 

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) 

420 

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] 

434 

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``. 

441 

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] 

464 

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)) 

469 

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 

477 

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()) 

496 

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) 

534 

535 def get_modified_times(self) -> dict[str, str]: 

536 """Return {document_id: stored_modified_isoformat} for all indexed documents. 

537 

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) 

544 

545 @property 

546 def _db_path(self) -> str: 

547 return str(Path(self._uri) / DB_FILENAME) 

548 

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. 

554 

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) 

572 

573 def compact(self, *, force: bool = False) -> None: 

574 """Rebuild the database file to reclaim space left behind by DELETEs. 

575 

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``. 

582 

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) 

611 

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) 

635 

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") 

664 

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) 

674 

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 

685 

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. 

690 

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 

699 

700 def check_and_run_migrations(self) -> bool: 

701 """Apply any pending schema migrations to the store. 

702 

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). 

707 

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 

716 

717 pending = sorted( 

718 [m for m in MIGRATIONS if current <= m.from_version < SCHEMA_VERSION], 

719 key=lambda m: m.from_version, 

720 ) 

721 

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) 

739 

740 return False 

741 

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) 

751 

752 

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