Coverage for src/lilbee/data/store/core.py: 100%

1145 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-09-28 17:20 +0000

1"""The ``Store`` class: high-level LanceDB read/write API used across lilbee.""" 

2 

3from __future__ import annotations 

4 

5import logging 

6import math 

7import os 

8from collections.abc import Callable, Iterable, Sequence 

9from contextlib import AbstractContextManager 

10from datetime import UTC, datetime 

11from pathlib import Path 

12from typing import TYPE_CHECKING, Final, Literal, cast 

13 

14import pyarrow as pa 

15 

16from lilbee.core.config import ( 

17 CHUNK_CONCEPTS_TABLE, 

18 CHUNKS_TABLE, 

19 CITATIONS_TABLE, 

20 ENTITIES_TABLE, 

21 ENTITY_SCHEMA_TABLE, 

22 INGEST_SOURCE_COLUMNS, 

23 MEMORIES_TABLE, 

24 META_TABLE, 

25 PAGE_TEXTS_TABLE, 

26 SOURCES_TABLE, 

27 WIKI_MENTIONS_TABLE, 

28 Config, 

29) 

30from lilbee.core.health_warnings import HealthWarning, WarningCode 

31from lilbee.core.vectors import Vector 

32from lilbee.retrieval.embedding_profiles import resolve_embedding_profile 

33from lilbee.runtime.lock import LOCK_TIMEOUT, LockTimeoutError, write_lock 

34 

35from .fusion import adaptive_weight_scale, fuse_arms, normalized_bm25, vector_similarity 

36from .lance_helpers import ( 

37 _CHUNK_COLUMN, 

38 _chunk_type_predicate, 

39 _has_fts_index, 

40 _has_scalar_index, 

41 _has_vector_index, 

42 _index_registry, 

43 _safe_delete_unlocked, 

44 _scalar_index_dangling, 

45 _sources_search_filter, 

46 _vector_index_dangling, 

47 ensure_table, 

48 escape_sql_string, 

49 refs_compatible, 

50 table_names, 

51) 

52from .ranking import mmr_rerank 

53from .schema import ( 

54 _citations_schema, 

55 _entity_schema_state_schema, 

56 _meta_schema, 

57 _page_texts_schema, 

58 _sources_schema, 

59 _wiki_mentions_schema, 

60) 

61from .types import ( 

62 ENTITY_SCHEMA_DELETE_ALL_PREDICATE, 

63 META_DELETE_ALL_PREDICATE, 

64 META_SCHEMA_VERSION, 

65 READ_CONSISTENCY_INTERVAL, 

66 SOURCE_STAT_UNKNOWN, 

67 ChunkType, 

68 ChunkWrite, 

69 CitationRecord, 

70 EmbeddingModelMismatchError, 

71 EntitySchemaState, 

72 MemoryKind, 

73 MemoryRow, 

74 PageTextRecord, 

75 RemoveResult, 

76 SearchChunk, 

77 SourceMeta, 

78 SourceRecord, 

79 SourceStat, 

80 SourceStatBackfill, 

81 SourceType, 

82 StoreMeta, 

83) 

84 

85if TYPE_CHECKING: 

86 import lance 

87 from lancedb.db import LanceDBConnection 

88 from lancedb.index import FTS, IndexConfig 

89 from lancedb.query import LanceFtsQueryBuilder 

90 from lancedb.table import LanceTable 

91 

92log = logging.getLogger(__name__) 

93 

94# Batched ingest flushes contend with long store operations (a search-triggered 

95# FTS optimize can hold the lock past the interactive 30s), and failing the 

96# flush replans and re-embeds the whole batch. Give the batch path more 

97# patience than interactive writes before it gives up. 

98BATCH_LOCK_TIMEOUT = 120.0 

99# Lock budget for index builds reached from the read path: a query must not 

100# stall behind a long ingest, so it skips the build and retries next search. 

101_READ_LOCK_TIMEOUT = 2.0 

102 

103 

104def _drop_unsupported_far_rows( 

105 results: list[SearchChunk], max_distance: float 

106) -> list[SearchChunk]: 

107 """Apply ``max_distance`` to rows whose only signal is the vector arm. 

108 

109 A row the BM25 arm also matched keeps lexical support regardless of its 

110 vector distance; dropping it on distance alone would re-bury exactly the 

111 identifier hits rank fusion exists to preserve. 

112 """ 

113 if max_distance <= 0: 

114 return results 

115 return [ 

116 r 

117 for r in results 

118 if r.bm25_score is not None or r.distance is None or r.distance <= max_distance 

119 ] 

120 

121 

122_MAX_THRESHOLD = 1.0 

123_MAX_FILTER_ITERATIONS = 20 # safety cap to prevent runaway loops 

124 

125 

126def _is_fts_position_overflow(exc: Exception) -> bool: 

127 """True when *exc* is LanceDB's positional-FTS list-encoding overflow. 

128 

129 A positional index (built by an intermediate dev commit) raises e.g. 

130 "Max offset N exceeds length of values M" on optimize(); a positionless 

131 rebuild is the remediation. Matched on message because LanceDB raises it 

132 as a generic error type. 

133 """ 

134 msg = str(exc).lower() 

135 return "offset" in msg and "exceeds" in msg 

136 

137 

138def _lexical_rows( 

139 table: LanceTable, 

140 query_text: str, 

141 limit: int, 

142 chunk_type: ChunkType | None, 

143 column: str = _CHUNK_COLUMN, 

144) -> list[SearchChunk]: 

145 """BM25 rows for *query_text* over a single FTS *column*. 

146 

147 ``MatchQuery`` pins the column and matches plain terms, so an unpinned search 

148 cannot widen to the title index and a quoted span cannot reach LanceDB as a 

149 phrase (which the positionless index rejects). This is the one place FTS 

150 queries are built; every arm goes through it. 

151 """ 

152 from lancedb.query import MatchQuery 

153 

154 # lancedb's stubs omit FullTextQuery from the fts overload and read MatchQuery's 

155 # kw-only defaults as required; the call is the documented FTS query form. 

156 query: LanceFtsQueryBuilder = table.search( # type: ignore[call-overload] 

157 MatchQuery(query_text, column), # type: ignore[call-arg] 

158 query_type="fts", 

159 ).limit(limit) 

160 if chunk_type: 

161 query = query.where(_chunk_type_predicate(chunk_type)) 

162 return [SearchChunk(**r) for r in query.to_list()] 

163 

164 

165# Vector ANN index. IVF_PQ compresses vectors so search scales to millions; 

166# refine_factor re-ranks the PQ candidates against full vectors to recover recall. 

167# The index type is carried by the lancedb IvfPq config at build time. 

168_VECTOR_METRIC: Final = "cosine" 

169# The scalar index kinds lilbee builds on the columns its queries filter by. 

170ScalarIndexType = Literal["BTREE", "BITMAP"] 

171_ANN_NPROBES_FLOOR = 20 

172# Fraction of IVF partitions probed per query. 0.05 was the "fast" end of the 

173# recall/latency curve and measurably cost recall at scale: on the 8.8M-passage 

174# MS MARCO index it gave recall@100 of 67%, where probing more partitions 

175# reached 87% (docs/benchmarks/retrieval-msmarco.md). 0.15 is the "balanced / 

176# high-recall" range for IVF and recovers most of that gap; refine_factor below 

177# still re-ranks the survivors against full vectors. 

178_ANN_NPROBES_PARTITION_FRACTION = 0.15 

179_ANN_REFINE_FACTOR = 10 

180 

181# Stat columns of ``_sources``; mirrors the field names in ``schema._sources_schema`` 

182# and ``types.SourceRecord``. Legacy tables that predate these columns are migrated 

183# in place with the SOURCE_STAT_UNKNOWN sentinel. 

184_SOURCE_STAT_COLUMNS = ("size_bytes", "mtime_ns", "stat_captured_ns") 

185 

186# Extraction-metadata columns of ``_sources``; nullable strings, so legacy 

187# tables migrate in place with NULL (meaning "extractor reported nothing"). 

188_SOURCE_META_COLUMNS = ("title", "authors", "created_at") 

189 

190# Document-title column of the chunks table; nullable so pre-title rows and 

191# writers that carry no title (wiki pages) read as NULL. 

192_TITLE_COLUMN = "title" 

193 

194# (table, source column) pairs deleted when a source's rows are replaced. The 

195# concept nodes/edges tables carry no source column (corpus-level aggregates), 

196# so only the per-chunk concept mapping is source-scoped. 

197_PER_SOURCE_TABLES = ( 

198 *( 

199 (name, INGEST_SOURCE_COLUMNS[name]) 

200 for name in (CHUNKS_TABLE, PAGE_TEXTS_TABLE, CHUNK_CONCEPTS_TABLE, ENTITIES_TABLE) 

201 ), 

202 # Per-(subject, source), so removing a source drops its mention evidence 

203 # here with its chunks; the wiki refresh only ever re-adds a source, never 

204 # has to remember to subtract a deleted one. Not in INGEST_SOURCE_COLUMNS 

205 # (a wiki table, not an ingest one), so its column is named directly. 

206 (WIKI_MENTIONS_TABLE, "source"), 

207) 

208 

209# (table, source column) pairs re-keyed when a source is relocated (moved on 

210# disk, same content). Extends the per-source set with the wiki citation's raw 

211# source_filename so citations keep pointing at the source after a move. 

212_RELOCATABLE_TABLES = ( 

213 *_PER_SOURCE_TABLES, 

214 (CITATIONS_TABLE, INGEST_SOURCE_COLUMNS[CITATIONS_TABLE]), 

215) 

216 

217# Sentinel: relocation must leave the stored title untouched (extraction-derived). 

218_KEEP_TITLE = "\x00keep" 

219 

220# Stat backfills replace this many source rows per locked write: the first 

221# sync after a stat-column upgrade backfills every source, and an unchunked 

222# replace would join millions of filenames into one delete predicate. 

223_SOURCE_STAT_BATCH_ROWS = 2000 

224 

225# Rows per Arrow batch when the aggregate scan walks the whole chunks table; 

226# bounds the decoded-text working set while the scan stays columnar. 

227_TERM_SCAN_BATCH_ROWS = 20_000 

228 

229# The title arm collapses each matched document to one row, so it over-fetches 

230# to gather enough distinct documents before deduping. Bounded so a title that 

231# hits a huge document can't scan the whole corpus. 

232_TITLE_FETCH_FACTOR = 20 

233_TITLE_MIN_FETCH = 200 

234_TITLE_FETCH_CEILING = 4096 

235 

236 

237def _ann_nprobes(row_count: int) -> int: 

238 """Partitions to probe: a fixed fraction of the IVF partition count (~sqrt(N)), floored.""" 

239 partitions = math.isqrt(max(row_count, 0)) 

240 return max(_ANN_NPROBES_FLOOR, math.ceil(partitions * _ANN_NPROBES_PARTITION_FRACTION)) 

241 

242 

243def _check_vector_dims(records: list[dict], embedding_dim: int) -> None: 

244 """Raise ``ValueError`` when any record's vector is not *embedding_dim* wide.""" 

245 for rec in records: 

246 vec = rec.get("vector", []) 

247 if len(vec) != embedding_dim: 

248 raise ValueError( 

249 f"Vector dimension mismatch: expected {embedding_dim}, " 

250 f"got {len(vec)} (source={rec.get('source', '?')})" 

251 ) 

252 

253 

254def _citations_for_wiki_predicate(wiki_source: str) -> str: 

255 """SQL predicate selecting every citation row belonging to *wiki_source*.""" 

256 return f"wiki_source = '{escape_sql_string(wiki_source)}'" 

257 

258 

259def _get_distance(chunk: SearchChunk) -> float: 

260 """Extract distance as a sortable float (inf for None).""" 

261 return chunk.distance if chunk.distance is not None else float("inf") 

262 

263 

264def _count_within_threshold(sorted_results: list[SearchChunk], threshold: float) -> int: 

265 """Count results whose distance is within the given threshold.""" 

266 for i, r in enumerate(sorted_results): 

267 if _get_distance(r) > threshold: 

268 return i 

269 return len(sorted_results) 

270 

271 

272class Store: 

273 """LanceDB vector store: wraps all DB operations with config-driven defaults.""" 

274 

275 def __init__(self, config: Config) -> None: 

276 self._config = config 

277 self._fts_ready: bool = False 

278 self._title_fts_ready: bool = False 

279 self._doc_prefix_warned: bool = False 

280 # Degradations the health endpoint reports; set where they are detected. 

281 self._fts_degraded: bool = False 

282 self._doc_prefix_mismatch: bool = False 

283 # Scalar indexes (source/chunk_type) are built at ingest; a serve-only 

284 # store builds them lazily from the search path. 

285 self._scalar_ready: bool = False 

286 self._db: LanceDBConnection | None = None 

287 # Cache of {filename: ingested_at} rebuilt only when sources 

288 # mutate; callers (temporal filter) hit it per-query. 

289 self._source_ingested_cache: dict[str, str] | None = None 

290 

291 def _fts_config(self) -> FTS: 

292 """Shared FTS index config: positionless (with_position=True overflows 

293 LanceDB's list encoding on optimize()) and stemmed for the configured 

294 corpus language.""" 

295 from lancedb.index import FTS 

296 

297 return FTS(with_position=False, language=self._config.fts_language) 

298 

299 def _index_build_lock(self, blocking: bool) -> AbstractContextManager[None]: 

300 """Write lock for index builds: full budget from ingest, short from the read path.""" 

301 return self._write_lock() if blocking else self._write_lock(_READ_LOCK_TIMEOUT) 

302 

303 def _write_lock(self, timeout: float = LOCK_TIMEOUT) -> AbstractContextManager[None]: 

304 """Acquire the write lock keyed on *this* store's data directory. 

305 

306 A per-instance ``Lilbee`` writes to its own ``lancedb_dir``; locking the 

307 global ``cfg`` dir instead would leave those writes uncoordinated across 

308 processes. 

309 """ 

310 return write_lock(self._config.lancedb_dir, timeout) 

311 

312 def _invalidate_source_cache(self) -> None: 

313 """Drop the cached {filename: ingested_at} map.""" 

314 self._source_ingested_cache = None 

315 

316 def source_ingested_at_map(self) -> dict[str, str]: 

317 """Return {filename: ingested_at} for every source, cached until mutation. 

318 

319 Best-effort: a reader racing a concurrent invalidation can store a 

320 pre-mutation snapshot. The only consumer (temporal query filter) 

321 treats a missing/stale entry as "do not filter," so staleness 

322 degrades ranking precision, never correctness. 

323 """ 

324 if self._source_ingested_cache is not None: 

325 return self._source_ingested_cache 

326 mapping = {s["filename"]: s.get("ingested_at", "") for s in self.get_sources()} 

327 self._source_ingested_cache = mapping 

328 return mapping 

329 

330 def _chunks_schema(self) -> pa.Schema: 

331 return pa.schema( 

332 [ 

333 pa.field("source", pa.utf8()), 

334 pa.field("content_type", pa.utf8()), 

335 pa.field("chunk_type", pa.utf8()), 

336 pa.field("page_start", pa.int32()), 

337 pa.field("page_end", pa.int32()), 

338 pa.field("line_start", pa.int32()), 

339 pa.field("line_end", pa.int32()), 

340 pa.field("chunk", pa.utf8()), 

341 pa.field("chunk_index", pa.int32()), 

342 pa.field(_TITLE_COLUMN, pa.utf8()), 

343 pa.field("vector", pa.list_(pa.float32(), self._config.embedding_dim)), 

344 ] 

345 ) 

346 

347 def _chunks_table(self) -> LanceTable: 

348 """Open/create the chunks table, adding the title column to pre-title tables.""" 

349 table = ensure_table(self.get_db(), CHUNKS_TABLE, self._chunks_schema()) 

350 if _TITLE_COLUMN not in table.schema.names: 

351 table.add_columns({_TITLE_COLUMN: "CAST(NULL AS STRING)"}) 

352 self._backfill_stem_titles_unlocked(table) 

353 return table 

354 

355 def _backfill_stem_titles_unlocked(self, table: LanceTable) -> None: 

356 """Backfill filename-stem titles for pre-upgrade rows. Caller holds ``write_lock()``. 

357 

358 Without this the title arm only matches documents ingested after the 

359 upgrade. Extracted (H1/EXIF) titles still need ``lilbee rebuild``; 

360 failure leaves NULLs, the pre-backfill behavior. 

361 """ 

362 from lilbee.data.title import derive_title # circular at module scope 

363 

364 try: 

365 rows = table.search().select(["source"]).limit(None).to_list() 

366 sources = sorted({r["source"] for r in rows}) 

367 filled = 0 

368 for source in sources: 

369 title = derive_title(source) 

370 if not title: 

371 continue 

372 escaped = source.replace("'", "''") 

373 table.update(where=f"source = '{escaped}'", values={_TITLE_COLUMN: title}) 

374 filled += 1 

375 log.info( 

376 "Backfilled filename titles for %d of %d existing sources; run " 

377 "`lilbee rebuild` to derive titles from document content", 

378 filled, 

379 len(sources), 

380 ) 

381 except Exception: 

382 log.warning( 

383 "Title backfill failed; pre-upgrade rows keep NULL titles until `lilbee rebuild`", 

384 exc_info=True, 

385 ) 

386 

387 def get_meta(self) -> StoreMeta | None: 

388 """Return the persisted store metadata row, or ``None`` if unset.""" 

389 table = self.open_table(META_TABLE) 

390 if table is None: 

391 return None 

392 rows = table.search().limit(None).to_list() 

393 if not rows: 

394 return None 

395 # _meta is meant to hold one row, but a swallowed delete on rewrite could 

396 # leave a stale one behind; take the newest so identity reads stay 

397 # deterministic rather than returning an arbitrary row. 

398 row = max(rows, key=lambda r: r["updated_at"]) 

399 return StoreMeta( 

400 embedding_model=row["embedding_model"], 

401 embedding_dim=int(row["embedding_dim"]), 

402 schema_version=int(row["schema_version"]), 

403 updated_at=row["updated_at"], 

404 ) 

405 

406 def _write_meta_unlocked(self, *, embedding_model: str, embedding_dim: int) -> None: 

407 """Overwrite the single ``_meta`` row with the supplied identity. 

408 

409 Caller must hold ``write_lock()``. Args are passed explicitly rather than 

410 re-read from ``self._config`` so the caller can snapshot cfg at a coherent 

411 instant and not race with a concurrent ``set_embedding_model``. 

412 """ 

413 db = self.get_db() 

414 table = ensure_table(db, META_TABLE, _meta_schema()) 

415 _safe_delete_unlocked(table, META_DELETE_ALL_PREDICATE) 

416 table.add( 

417 [ 

418 { 

419 "embedding_model": embedding_model, 

420 "embedding_dim": embedding_dim, 

421 "schema_version": META_SCHEMA_VERSION, 

422 "updated_at": datetime.now(UTC).isoformat(), 

423 } 

424 ] 

425 ) 

426 # A rebuild is the remedy the prefix warning names, so the next query 

427 # re-decides instead of logging a mismatch this write just resolved. 

428 self._doc_prefix_warned = False 

429 

430 def _has_chunks(self) -> bool: 

431 """Return True when the chunks table exists and has at least one row.""" 

432 chunks = self.open_table(CHUNKS_TABLE) 

433 return chunks is not None and chunks.count_rows() > 0 

434 

435 def has_chunks(self) -> bool: 

436 """Public predicate: True iff the store currently holds at least one chunk.""" 

437 return self._has_chunks() 

438 

439 def initialize_meta_if_legacy(self) -> bool: 

440 """Pin a legacy store's identity to the current cfg if not already set. 

441 

442 Returns ``True`` when a meta row was just written. No-op when meta already 

443 exists or no chunks are present. This is the path that converts a 

444 pre-upgrade store (chunks present, no ``_meta``) into a gated store. It 

445 snapshots cfg under the write lock to keep the recorded identity coherent 

446 with what the gate is comparing against. 

447 """ 

448 if self.get_meta() is not None: 

449 return False 

450 if not self._has_chunks(): 

451 return False 

452 embedding_model = self._config.embedding_model 

453 embedding_dim = self._config.embedding_dim 

454 with self._write_lock(): 

455 # Re-check under the lock so two callers do not both warn-and-write. 

456 if self.get_meta() is not None: 

457 return False 

458 log.warning( 

459 "Legacy store has chunks but no _meta row. Initializing _meta from " 

460 "the current configuration (embedding_model=%s, embedding_dim=%d). " 

461 "If you changed embedding_model before upgrading, run `lilbee rebuild` " 

462 "to ensure the store is consistent.", 

463 embedding_model, 

464 embedding_dim, 

465 ) 

466 self._write_meta_unlocked(embedding_model=embedding_model, embedding_dim=embedding_dim) 

467 return True 

468 

469 def index_mismatch(self) -> EmbeddingModelMismatchError | None: 

470 """The drift between the persisted embedding identity and cfg, or None. 

471 

472 Pure check, no side effects, and the one verdict search, ingest, the 

473 embedding setter, sync and health all read. Migration of legacy stores 

474 (chunks present, no ``_meta``) is the caller's responsibility via 

475 ``initialize_meta_if_legacy``; rewriting a legacy bare-repo ``_meta`` row 

476 to the canonical full ref is the caller's responsibility via 

477 ``canonicalize_meta_if_legacy``. Safe to call from inside an existing 

478 ``write_lock()`` (no recursive lock attempt). cfg fields are snapshotted 

479 at entry so the comparison is coherent even if another thread mutates 

480 them mid-call. 

481 """ 

482 current_model = self._config.embedding_model 

483 current_dim = self._config.embedding_dim 

484 meta = self.get_meta() 

485 if meta is None: 

486 return None 

487 if refs_compatible( 

488 meta["embedding_model"], current_model, meta["embedding_dim"], current_dim 

489 ): 

490 return None 

491 return EmbeddingModelMismatchError( 

492 persisted_model=meta["embedding_model"], 

493 persisted_dim=meta["embedding_dim"], 

494 current_model=current_model, 

495 current_dim=current_dim, 

496 ) 

497 

498 def _ensure_embedding_compat(self) -> None: 

499 """Raise when the persisted embedding identity drifts from cfg.""" 

500 mismatch = self.index_mismatch() 

501 if mismatch is not None: 

502 raise mismatch 

503 

504 def _doc_prefix_is_stale(self) -> bool: 

505 """Whether the stored documents predate the embedding family's document prefix.""" 

506 meta = self.get_meta() 

507 if meta is None: 

508 return False 

509 profile = resolve_embedding_profile(self._config.embedding_model) 

510 return bool(profile.doc_prefix) and meta["schema_version"] < profile.doc_prefix_since 

511 

512 def _warn_stale_doc_prefix(self) -> None: 

513 """Warn once when the embedding family's document prefix postdates this store. 

514 

515 Queries would carry the family prefix while stored documents do not; 

516 the mismatch degrades quality silently until a rebuild re-embeds. The 

517 check reads the metadata row, which the query path cannot pay per search. 

518 """ 

519 if self._doc_prefix_warned: 

520 return 

521 self._doc_prefix_warned = True 

522 if not self._doc_prefix_is_stale(): 

523 return 

524 self._doc_prefix_mismatch = True 

525 log.warning( 

526 "This index predates '%s' document prefixes: queries are prefixed " 

527 "but stored documents are not. Run `lilbee rebuild` to re-embed.", 

528 self._config.embedding_model, 

529 ) 

530 

531 def _index_mismatch_warning(self) -> HealthWarning | None: 

532 """The stale-index degradation, from the same verdict search refuses on.""" 

533 mismatch = self.index_mismatch() 

534 if mismatch is None: 

535 return None 

536 return HealthWarning( 

537 code=WarningCode.INDEX_EMBEDDING_MISMATCH, 

538 message=( 

539 f"The index was built with embedding model '{mismatch.persisted_model}' " 

540 f"({mismatch.persisted_dim} dims), but '{mismatch.current_model}' " 

541 f"({mismatch.current_dim} dims) is configured, so search refuses it." 

542 ), 

543 remedy=( 

544 f"Rebuild the index under '{mismatch.current_model}', or switch the " 

545 f"embedding model back to '{mismatch.persisted_model}'." 

546 ), 

547 ) 

548 

549 def health_warnings(self) -> list[HealthWarning]: 

550 """Retrieval degradations a client should know about before it searches.""" 

551 warnings: list[HealthWarning] = [] 

552 if (stale := self._index_mismatch_warning()) is not None: 

553 warnings.append(stale) 

554 if self._fts_degraded: 

555 warnings.append( 

556 HealthWarning( 

557 code=WarningCode.FTS_UNAVAILABLE, 

558 message=( 

559 "Keyword search is unavailable, so answers are drawn from " 

560 "vector similarity alone and exact terms may be missed." 

561 ), 

562 remedy="Run `lilbee rebuild` to rebuild the search index.", 

563 ) 

564 ) 

565 table = self.open_table(CHUNKS_TABLE) 

566 indices = _index_registry(table) if table is not None else None 

567 if table is not None and indices is not None: 

568 dangling_scalar = _scalar_index_dangling(table, indices, self._config.lancedb_dir) 

569 if dangling_scalar: 

570 warnings.append( 

571 HealthWarning( 

572 code=WarningCode.SCALAR_INDEX_UNAVAILABLE, 

573 message=( 

574 "The scalar index on " 

575 + ", ".join(dangling_scalar) 

576 + " is missing its files, so filtered search is degraded." 

577 ), 

578 remedy="Run `lilbee rebuild` to rebuild the search index.", 

579 ) 

580 ) 

581 if self._doc_prefix_mismatch and not self._doc_prefix_is_stale(): 

582 # The remedy is `lilbee rebuild`, which usually runs in another 

583 # process, so the query path's cached verdict outlives the fix. 

584 self._doc_prefix_mismatch = False 

585 if self._doc_prefix_mismatch: 

586 warnings.append( 

587 HealthWarning( 

588 code=WarningCode.EMBEDDING_PREFIX_MISMATCH, 

589 message=( 

590 f"This index predates '{self._config.embedding_model}' document " 

591 "prefixes: queries are prefixed but stored documents are not, " 

592 "so retrieval is less accurate." 

593 ), 

594 remedy="Run `lilbee rebuild` to re-embed the documents.", 

595 ) 

596 ) 

597 return warnings 

598 

599 def assert_embedding_compatible(self) -> None: 

600 """Run the full embedding-identity gate (legacy init, canonicalize, check). 

601 

602 Mirrors the gate ``search`` applies. Callers that write under a fresh 

603 embedder (import) use this to fail before any destructive work when the 

604 store was built by a different model. 

605 """ 

606 self.initialize_meta_if_legacy() 

607 self.canonicalize_meta_if_legacy() 

608 self._ensure_embedding_compat() 

609 

610 def _needs_canonical_meta_rewrite( 

611 self, meta: StoreMeta | None, current_model: str, current_dim: int 

612 ) -> bool: 

613 """True iff *meta* is the legacy form and refs-compatible with current cfg.""" 

614 if meta is None or meta["embedding_model"] == current_model: 

615 return False 

616 return refs_compatible( 

617 meta["embedding_model"], current_model, meta["embedding_dim"], current_dim 

618 ) 

619 

620 def canonicalize_meta_if_legacy(self) -> bool: 

621 """Rewrite a legacy bare-repo ``_meta`` row to the canonical full ref. 

622 

623 Pre-canonical lilbee persisted only ``<org>/<repo>`` in 

624 ``_meta.embedding_model``. The current code persists the full 

625 ``<org>/<repo>/<filename>.gguf``. When the two refer to the same 

626 model under :func:`refs_compatible` but differ as raw strings, the 

627 meta row is rewritten so the legacy name never surfaces. Returns 

628 ``True`` on write; ``False`` when missing, already canonical, or 

629 incompatible (the gate handles incompatibility). 

630 """ 

631 current_model = self._config.embedding_model 

632 current_dim = self._config.embedding_dim 

633 if not self._needs_canonical_meta_rewrite(self.get_meta(), current_model, current_dim): 

634 return False 

635 with self._write_lock(): 

636 meta = self.get_meta() # re-read under the lock for racing callers 

637 if not self._needs_canonical_meta_rewrite(meta, current_model, current_dim): 

638 return False 

639 assert meta is not None # filtered above # noqa: S101 

640 log.info( 

641 "Migrating legacy embedding ref in store meta: %r -> %r", 

642 meta["embedding_model"], 

643 current_model, 

644 ) 

645 self._write_meta_unlocked(embedding_model=current_model, embedding_dim=current_dim) 

646 return True 

647 

648 def get_db(self) -> LanceDBConnection: 

649 if self._db is None: 

650 from lancedb.db import LanceDBConnection 

651 

652 self._config.lancedb_dir.mkdir(parents=True, exist_ok=True) 

653 self._db = LanceDBConnection( 

654 str(self._config.lancedb_dir), 

655 read_consistency_interval=READ_CONSISTENCY_INTERVAL, 

656 ) 

657 return self._db 

658 

659 def open_table(self, name: str) -> LanceTable | None: 

660 """Open a table if it exists, otherwise return None.""" 

661 db = self.get_db() 

662 if name not in table_names(db): 

663 return None 

664 table: LanceTable = db.open_table(name) 

665 return table 

666 

667 def ensure_fts_index(self, *, blocking: bool = True) -> None: 

668 """Create the chunks FTS index, or run ``optimize()`` once it exists. 

669 

670 ``optimize()`` folds newly added rows into the FTS index and also 

671 runs LanceDB's default compaction + version pruning (default prune 

672 window: 7 days). Work scales with recent deltas rather than total 

673 chunk count, so large corpora no longer pay the full 

674 ``create_index(config=FTS(), replace=True)`` rebuild cost on every sync. 

675 

676 ``blocking=False`` (the search path) marks an existing index ready 

677 without the lock and skips maintenance when another process holds it, 

678 so a long concurrent ingest cannot stall or fail a query. An unreadable 

679 index registry never counts as an existing index. 

680 """ 

681 probe = self.open_table(CHUNKS_TABLE) 

682 if probe is None: 

683 return 

684 indices = _index_registry(probe) 

685 if indices is not None and _has_fts_index(indices): 

686 self._fts_ready = True 

687 try: 

688 with self._index_build_lock(blocking): 

689 self._ensure_fts_index_unlocked() 

690 except LockTimeoutError: 

691 if blocking: 

692 raise 

693 log.debug("Skipped FTS index maintenance; another process holds the write lock") 

694 

695 def _ensure_fts_index_unlocked(self) -> None: 

696 """Body of ``ensure_fts_index``. Caller holds ``write_lock()``.""" 

697 table = self.open_table(CHUNKS_TABLE) 

698 if table is None: 

699 return 

700 try: 

701 indices = _index_registry(table) 

702 if indices is None: 

703 self._repair_fts_registry(table) 

704 return 

705 if _has_fts_index(indices): 

706 self._fts_ready = True 

707 try: 

708 # One optimize folds new rows into every index on the table. 

709 table.optimize() 

710 log.debug("FTS index optimized on '%s'", CHUNKS_TABLE) 

711 if _index_registry(table) is None: 

712 # Suspected producer of an index without files; unconfirmed. 

713 self._repair_fts_registry(table) 

714 return 

715 except Exception as exc: 

716 if _is_fts_position_overflow(exc): 

717 # Positional indexes overflow on optimize(); rebuild 

718 # positionless once. 

719 self._rebuild_fts(table, "a positional index overflowed on optimize()") 

720 else: 

721 log.warning( 

722 "FTS optimize() failed; the existing index still serves hybrid search", 

723 exc_info=True, 

724 ) 

725 else: 

726 # Positionless: with_position=True overflows LanceDB's list 

727 # encoding on optimize(), and nothing issues phrase queries. 

728 table.create_index(_CHUNK_COLUMN, config=self._fts_config(), replace=False) 

729 self._fts_ready = True 

730 log.debug("FTS index created on '%s'", CHUNKS_TABLE) 

731 # Only the opt-in title arm needs the title index. 

732 if self._config.title_search: 

733 self._ensure_title_fts_unlocked(table) 

734 except Exception: 

735 log.debug("FTS index ensure failed (empty table?)", exc_info=True) 

736 

737 def _ensure_title_fts_unlocked(self, table: LanceTable) -> None: 

738 """Create the title FTS index when the column exists. Caller holds ``write_lock()``. 

739 

740 Failure never blocks the chunk index: the title arm feature-detects the 

741 index per query, so a store without it simply searches without titles. 

742 """ 

743 indices = _index_registry(table) 

744 if indices is None: 

745 self._repair_fts_registry(table) 

746 return 

747 if _TITLE_COLUMN not in table.schema.names or _has_fts_index(indices, _TITLE_COLUMN): 

748 self._title_fts_ready = _has_fts_index(indices, _TITLE_COLUMN) 

749 return 

750 try: 

751 # Positionless for the same reason as the chunk index. 

752 table.create_index(_TITLE_COLUMN, config=self._fts_config(), replace=False) 

753 self._title_fts_ready = True 

754 log.debug("Title FTS index created on '%s'", CHUNKS_TABLE) 

755 except Exception: 

756 # Only reached with title_search enabled, so a silent failure means 

757 # the user's opted-in title arm quietly does nothing. Warn, don't hide. 

758 log.warning( 

759 "Title FTS index creation failed; the title-search arm will " 

760 "contribute nothing until it can be built", 

761 exc_info=True, 

762 ) 

763 

764 def ensure_title_fts_index(self, *, blocking: bool = True) -> None: 

765 """Build the title FTS index for a title_search toggle after startup. 

766 

767 Without this, a process that latched ``_fts_ready`` before the toggle 

768 never builds the index and the title arm stays silently dead until the 

769 next ingest or restart. 

770 """ 

771 table = self.open_table(CHUNKS_TABLE) 

772 if table is None: 

773 return 

774 indices = _index_registry(table) 

775 if indices is not None and _has_fts_index(indices, _TITLE_COLUMN): 

776 self._title_fts_ready = True 

777 return 

778 try: 

779 with self._index_build_lock(blocking): 

780 self._ensure_title_fts_unlocked(table) 

781 except LockTimeoutError: 

782 if blocking: 

783 raise 

784 log.debug("Skipped title FTS build; another process holds the write lock") 

785 

786 def _rebuild_fts(self, table: LanceTable, reason: str) -> None: 

787 """Replace the FTS indexes with fresh positionless ones. Caller holds the lock. 

788 

789 The one-shot remediation for an index built ``with_position=True`` that 

790 overflows on every ``optimize()``. The title index is rebuilt too when 

791 the title arm is enabled. 

792 """ 

793 try: 

794 table.create_index(_CHUNK_COLUMN, config=self._fts_config(), replace=True) 

795 if self._config.title_search and _TITLE_COLUMN in table.schema.names: 

796 table.create_index(_TITLE_COLUMN, config=self._fts_config(), replace=True) 

797 log.warning("Rebuilt the FTS index because %s", reason) 

798 except Exception: 

799 log.warning( 

800 "Positionless FTS rebuild failed; the existing index still serves", 

801 exc_info=True, 

802 ) 

803 

804 def _repair_fts_registry(self, table: LanceTable) -> list[IndexConfig] | None: 

805 """Rebuild the chunks FTS indexes after the index listing raised. Caller holds the lock. 

806 

807 lancedb opens each FTS index's files while it lists the indexes, so an 

808 FTS index without its files makes the listing raise; scalar and vector 

809 indexes do not. Each FTS index (chunk text, plus title when the title 

810 arm is on) is rebuilt with ``replace=True`` in its own try. The listing 

811 is then read once more; the directory check on a readable listing 

812 decides whether a scalar or vector index needs its own rebuild. 

813 

814 Returns the listing after the rebuild, ``None`` (logged once) when it 

815 is still unreadable. 

816 """ 

817 columns = [_CHUNK_COLUMN] 

818 if self._config.title_search and _TITLE_COLUMN in table.schema.names: 

819 columns.append(_TITLE_COLUMN) 

820 for column in columns: 

821 try: 

822 table.create_index(column, config=self._fts_config(), replace=True) 

823 except Exception: 

824 log.warning( 

825 "Could not rebuild the FTS index on '%s.%s'", table.name, column, exc_info=True 

826 ) 

827 continue 

828 log.warning( 

829 "Rebuilt the FTS index on '%s.%s' because the index listing is unreadable", 

830 table.name, 

831 column, 

832 ) 

833 if column == _CHUNK_COLUMN: 

834 self._fts_ready = True 

835 else: 

836 self._title_fts_ready = True 

837 indices = _index_registry(table) 

838 if indices is None: 

839 log.warning( 

840 "The index listing on '%s' is still unreadable after the FTS rebuild; " 

841 "run `lilbee rebuild` to rebuild the search index", 

842 table.name, 

843 ) 

844 return indices 

845 

846 # Tables and (column, index_type) pairs the query path filters by. 

847 # chunk_concepts serves the concept boost (ConceptGraph._chunk_concepts_from); 

848 # without its index every boosted query full-scans the table. 

849 _SCALAR_TARGETS: tuple[tuple[str, tuple[tuple[str, ScalarIndexType], ...]], ...] = ( 

850 (CHUNKS_TABLE, (("source", "BTREE"), ("chunk_type", "BITMAP"))), 

851 (CHUNK_CONCEPTS_TABLE, (("chunk_source", "BTREE"),)), 

852 ) 

853 

854 def ensure_scalar_indexes(self, *, blocking: bool = True) -> None: 

855 """Build scalar indexes on the columns lilbee filters by. 

856 

857 ``source`` and ``chunk_type`` predicates run as prefilters (LanceDB's 

858 default), but without an index each is a full-table scan. Readiness 

859 latches only when every target table exists and is covered, so a table 

860 created later (chunk_concepts under serve ordering) still gets its 

861 index on a following call. A table whose scalar index lost its files 

862 counts as pending, so the build step replaces it. A table whose index 

863 listing is unreadable keeps readiness unlatched until the FTS path 

864 repairs the listing. The lock is taken only 

865 when there is something to build; ``blocking=False`` (the search path) 

866 skips the build when another process holds it instead of stalling the 

867 query. 

868 """ 

869 pending = [] 

870 complete = True 

871 for name, columns in self._SCALAR_TARGETS: 

872 table = self.open_table(name) 

873 if table is None: 

874 complete = False 

875 continue 

876 indices = _index_registry(table) 

877 if indices is None: 

878 # The FTS path repairs an unreadable listing; check again then. 

879 complete = False 

880 continue 

881 names = table.schema.names 

882 dangling = _scalar_index_dangling(table, indices, self._config.lancedb_dir) 

883 needs_build = any( 

884 c in names and (not _has_scalar_index(indices, c) or c in dangling) 

885 for c, _ in columns 

886 ) 

887 if needs_build: 

888 pending.append((name, columns)) 

889 if not pending: 

890 self._scalar_ready = complete 

891 return 

892 try: 

893 with self._index_build_lock(blocking): 

894 for name, columns in pending: 

895 self._ensure_scalar_index_on(name, columns) 

896 self._scalar_ready = complete 

897 except LockTimeoutError: 

898 if blocking: 

899 raise 

900 log.debug("Skipped scalar index build; another process holds the write lock") 

901 

902 def _ensure_scalar_index_on( 

903 self, table_name: str, columns: tuple[tuple[str, ScalarIndexType], ...] 

904 ) -> None: 

905 """Build the given (column, index_type) scalar indexes on *table_name*. 

906 

907 Caller holds ``write_lock()``. A registered index whose files are gone 

908 is replaced; nothing is built while the index listing is unreadable. 

909 Each column gets its own try so one failure does not skip the rest; a 

910 failure on a populated table warns (the prefilter speedup is silently 

911 lost) while an empty table's is debug. 

912 """ 

913 table = self.open_table(table_name) 

914 if table is None: 

915 return 

916 indices = _index_registry(table) 

917 if indices is None: 

918 return 

919 names = table.schema.names 

920 fail_level = logging.WARNING if table.count_rows() > 0 else logging.DEBUG 

921 dangling = _scalar_index_dangling(table, indices, self._config.lancedb_dir) 

922 for column, index_type in columns: 

923 if column not in names or ( 

924 _has_scalar_index(indices, column) and column not in dangling 

925 ): 

926 continue 

927 try: 

928 table.create_scalar_index(column, index_type=index_type, replace=column in dangling) 

929 log.debug("Scalar (%s) index created on '%s.%s'", index_type, table_name, column) 

930 except Exception: 

931 log.log( 

932 fail_level, 

933 "Scalar index create failed on '%s.%s'", 

934 table_name, 

935 column, 

936 exc_info=True, 

937 ) 

938 

939 def _create_vector_index(self, table: LanceTable) -> None: 

940 """Build the IVF_PQ index on the vector column, replacing any existing one.""" 

941 from lancedb.index import IvfPq 

942 

943 table.create_index("vector", config=IvfPq(distance_type=_VECTOR_METRIC), replace=True) 

944 

945 def ensure_vector_index(self, *, force: bool = False) -> bool: 

946 """Build or refresh the ANN vector index when the corpus is large enough. 

947 

948 Below ``cfg.ann_index_threshold`` (or when it is 0) the store keeps exact 

949 flat search, which is faster and exact for small vaults and is all a 

950 laptop needs. Once an index exists, ``optimize()`` folds new rows in. 

951 Pass ``force=True`` to build regardless of the threshold (publish flow). 

952 An unreadable index listing gets its FTS repair first, and a registered 

953 index whose files are gone is rebuilt instead of optimized. Returns 

954 True when an index was created, refreshed, or rebuilt. 

955 """ 

956 threshold = self._config.ann_index_threshold 

957 with self._write_lock(): 

958 table = self.open_table(CHUNKS_TABLE) 

959 if table is None: 

960 return False 

961 indices = _index_registry(table) 

962 if indices is None: 

963 indices = self._repair_fts_registry(table) 

964 if indices is None: 

965 return False 

966 if _vector_index_dangling(table, indices, self._config.lancedb_dir): 

967 return self._rebuild_dangling_vector_index(table) 

968 if _has_vector_index(indices): 

969 table.optimize() 

970 log.debug("Vector index optimized on '%s'", CHUNKS_TABLE) 

971 if _index_registry(table) is None: 

972 # Repair in the same step when the prune dropped FTS files. 

973 self._repair_fts_registry(table) 

974 return True 

975 if not force and (threshold <= 0 or table.count_rows() < threshold): 

976 return False 

977 return self._build_vector_index(table) 

978 

979 def _rebuild_dangling_vector_index(self, table: LanceTable) -> bool: 

980 """Replace a registered vector index whose files are gone. Caller holds the lock.""" 

981 try: 

982 self._create_vector_index(table) 

983 except Exception: 

984 log.warning( 

985 "Could not rebuild the vector index on '%s' whose files are missing", 

986 CHUNKS_TABLE, 

987 exc_info=True, 

988 ) 

989 return False 

990 log.warning("Rebuilt the vector index on '%s' because its files are missing", CHUNKS_TABLE) 

991 return True 

992 

993 def _build_vector_index(self, table: LanceTable) -> bool: 

994 """Build the first vector index, warning on failure. Caller holds the lock.""" 

995 try: 

996 self._create_vector_index(table) 

997 log.info("Vector ANN index created on '%s'", CHUNKS_TABLE) 

998 return True 

999 except Exception: 

1000 log.warning( 

1001 "Vector ANN index build failed on '%s' at %d rows; search falls back " 

1002 "to exact flat scan, which is slow at this scale. Free up memory/disk " 

1003 "and re-run to rebuild the index.", 

1004 CHUNKS_TABLE, 

1005 table.count_rows(), 

1006 exc_info=True, 

1007 ) 

1008 return False 

1009 

1010 def _add_chunks_unlocked(self, records: list[dict]) -> int: 

1011 """Add chunk records and return the count. Caller must hold ``write_lock()``.""" 

1012 embedding_model = self._config.embedding_model 

1013 embedding_dim = self._config.embedding_dim 

1014 self._ensure_embedding_compat() 

1015 self._fts_ready = False 

1016 self._scalar_ready = False 

1017 if not records: 

1018 return 0 

1019 _check_vector_dims(records, embedding_dim) 

1020 table = self._chunks_table() 

1021 table.add(records) 

1022 if self.get_meta() is None: 

1023 self._write_meta_unlocked(embedding_model=embedding_model, embedding_dim=embedding_dim) 

1024 return len(records) 

1025 

1026 def add_chunks(self, records: list[dict]) -> int: 

1027 """Add chunk records to the store. Returns count added. 

1028 

1029 Raises ``EmbeddingModelMismatchError`` if the persisted ``_meta`` row was 

1030 written under a different embedding model than the current ``cfg``. On the 

1031 first write to a fresh store, ``_meta`` is initialized from the current cfg. 

1032 

1033 The gate runs inside the write lock and uses a single cfg snapshot so a 

1034 concurrent ``set_embedding_model`` cannot slip a write in past a stale 

1035 compatibility check. 

1036 """ 

1037 with self._write_lock(): 

1038 return self._add_chunks_unlocked(records) 

1039 

1040 def replace_chunks(self, records: list[dict], predicate: str) -> int: 

1041 """Replace the chunk rows matching *predicate* with *records* under one write lock. 

1042 

1043 Same compatibility and dimension gates as :meth:`add_chunks`, both run 

1044 before the delete so a rejected write leaves the existing rows in place. 

1045 The lock serializes writers; delete and add are separate commits, so a 

1046 concurrent reader can briefly see the rows absent, and callers retry on 

1047 a crash between the two. A delete failure propagates, following the same 

1048 rule as :meth:`_delete_by_sources_unlocked`: swallowed, the caller would 

1049 read a success-shaped result over rows still describing the old body. 

1050 Returns the count added. 

1051 """ 

1052 with self._write_lock(): 

1053 self._ensure_embedding_compat() 

1054 _check_vector_dims(records, self._config.embedding_dim) 

1055 table = self._chunks_table() 

1056 table.delete(predicate) 

1057 return self._add_chunks_unlocked(records) 

1058 

1059 def _stamp_meta_unlocked(self, embedding_model: str, embedding_dim: int) -> None: 

1060 """Write the embedder identity on the first write to a fresh store.""" 

1061 if self.get_meta() is None: 

1062 self._write_meta_unlocked(embedding_model=embedding_model, embedding_dim=embedding_dim) 

1063 

1064 def absorb_rows(self, name: str, rows: pa.Table) -> int: 

1065 """Append an ingest shard's *rows* to table *name*, creating it if absent. 

1066 

1067 The rows already carry their embeddings, so this is the merge path for a 

1068 multi-GPU sync: the per-worker stores are folded in whole and the indexes 

1069 are rebuilt corpus-wide afterwards. 

1070 """ 

1071 with self._write_lock(): 

1072 self._ensure_embedding_compat() 

1073 self._fts_ready = False 

1074 self._scalar_ready = False 

1075 ensure_table(self.get_db(), name, rows.schema).add(rows) 

1076 self._stamp_meta_unlocked(self._config.embedding_model, self._config.embedding_dim) 

1077 return int(rows.num_rows) 

1078 

1079 def _assert_schemas_match( 

1080 self, 

1081 name: str, 

1082 target: Path, 

1083 shard_tables: list[Path], 

1084 sources: Sequence[lance.LanceDataset], 

1085 ) -> None: 

1086 """Refuse shards whose schema differs from the table's. 

1087 

1088 A mismatched vector width is the realistic case: two workers built with 

1089 different embedding models. Adoption commits fragment metadata and never 

1090 passes the rows through a writer, so nothing else would notice. 

1091 """ 

1092 import lance 

1093 

1094 expected = lance.dataset(str(target)).schema 

1095 for shard_table, source in zip(shard_tables, sources, strict=True): 

1096 if not source.schema.equals(expected): 

1097 raise ValueError( 

1098 f"Cannot adopt {shard_table} into {name}: its schema does not match " 

1099 f"the index. A shard built with a different embedding model cannot be " 

1100 f"folded in; re-ingest it with the model this index was built on." 

1101 ) 

1102 

1103 def adopt_fragments(self, name: str, shard_tables: list[Path]) -> int: 

1104 """Take over *shard_tables*' data files for table *name* without copying rows. 

1105 

1106 Every fragment's data file is hard-linked into this table's directory and 

1107 the whole set committed as one metadata-only append, so the rows are never 

1108 read or rewritten and the bytes exist once on disk with two names. That is 

1109 the difference between a merge that costs the corpus and one that costs its 

1110 fragment count: the copy path rewrites every vector, and because the shard 

1111 stores are kept as resume state the corpus would then be on disk twice. 

1112 

1113 Whole-fragment only, so this serves the first full merge. A re-sync merges 

1114 named sources, where a fragment holds both touched and untouched rows and 

1115 the scoped row copy is both correct and already cheap. 

1116 

1117 Returns the rows adopted. Raises ValueError when a shard's schema differs 

1118 from the table's, and OSError when a shard is on another filesystem (hard 

1119 links cannot cross one) or a data file name collides. The parent is left 

1120 as it was in every failure: schemas are checked before anything is linked, 

1121 links made before a failure are removed, and the commit happens once at 

1122 the end. 

1123 """ 

1124 import lance 

1125 

1126 with self._write_lock(): 

1127 self._ensure_embedding_compat() 

1128 target = self._config.lancedb_dir / f"{name}.lance" 

1129 sources = [lance.dataset(str(shard_table)) for shard_table in shard_tables] 

1130 if not sources: 

1131 return 0 

1132 if name not in table_names(self.get_db()): 

1133 ensure_table(self.get_db(), name, sources[0].schema) 

1134 # Before anything is linked. Committing a fragment whose schema does 

1135 # not match writes an index that reads back as a panic inside Arrow 

1136 # rather than an error: the row copy is rejected by the writer, and 

1137 # adoption has no writer to reject it. 

1138 self._assert_schemas_match(name, target, shard_tables, sources) 

1139 # A table created empty has no data directory yet: nothing has been 

1140 # written into it, and the links below need somewhere to go. 

1141 (target / "data").mkdir(parents=True, exist_ok=True) 

1142 adopted: list[lance.FragmentMetadata] = [] 

1143 linked: list[Path] = [] 

1144 rows = 0 

1145 try: 

1146 for shard_table, source in zip(shard_tables, sources, strict=True): 

1147 for fragment in source.get_fragments(): 

1148 meta = fragment.metadata 

1149 for data_file in meta.files: 

1150 filename = Path(data_file.path).name 

1151 link = target / "data" / filename 

1152 os.link(shard_table / "data" / filename, link) 

1153 linked.append(link) 

1154 adopted.append(meta) 

1155 rows += source.count_rows() 

1156 except OSError: 

1157 # A link made before the failure is a file the manifest never 

1158 # names, so nothing would ever remove it. The caller falls back 

1159 # to copying rows. 

1160 for link in linked: 

1161 link.unlink(missing_ok=True) 

1162 raise 

1163 if not adopted: 

1164 return 0 

1165 self._fts_ready = False 

1166 self._scalar_ready = False 

1167 existing = lance.dataset(str(target)) 

1168 lance.LanceDataset.commit( 

1169 str(target), 

1170 lance.LanceOperation.Append(adopted), 

1171 read_version=existing.version, 

1172 ) 

1173 self._stamp_meta_unlocked(self._config.embedding_model, self._config.embedding_dim) 

1174 return rows 

1175 

1176 def bm25_probe( 

1177 self, query_text: str, top_k: int = 5, chunk_type: ChunkType | None = None 

1178 ) -> list[SearchChunk]: 

1179 """Quick BM25-only search for confidence checking. Returns up to top_k results. 

1180 

1181 When *chunk_type* is set, only chunks of that type ("raw" or "wiki") are returned. 

1182 """ 

1183 table = self.open_table(CHUNKS_TABLE) 

1184 if table is None: 

1185 return [] 

1186 if not self._fts_ready: 

1187 self.ensure_fts_index(blocking=False) 

1188 if not self._fts_ready: 

1189 return [] 

1190 try: 

1191 results = _lexical_rows(table, query_text, top_k, chunk_type) 

1192 norms = normalized_bm25([r.bm25_score or 0.0 for r in results]) 

1193 return [ 

1194 r.model_copy(update={"score": norm}) for r, norm in zip(results, norms, strict=True) 

1195 ] 

1196 except Exception: 

1197 log.debug("BM25 probe failed", exc_info=True) 

1198 return [] 

1199 

1200 def search( 

1201 self, 

1202 query_vector: Vector, 

1203 top_k: int | None = None, 

1204 max_distance: float | None = None, 

1205 query_text: str | None = None, 

1206 chunk_type: ChunkType | None = None, 

1207 ) -> list[SearchChunk]: 

1208 """Search for similar chunks. Hybrid when FTS available, else vector-only. 

1209 

1210 Results with distance > max_distance are filtered out (vector-only path). 

1211 Pass max_distance=0 to disable filtering. 

1212 When *chunk_type* is set, only chunks of that type ("raw" or "wiki") are returned. 

1213 

1214 Raises ``EmbeddingModelMismatchError`` if the persisted ``_meta`` row was 

1215 written under a different embedding model than the current ``cfg``. 

1216 """ 

1217 if top_k is None: 

1218 top_k = self._config.top_k 

1219 if max_distance is None: 

1220 max_distance = self._config.max_distance 

1221 table = self.open_table(CHUNKS_TABLE) 

1222 if table is None: 

1223 return [] 

1224 self.initialize_meta_if_legacy() 

1225 self.canonicalize_meta_if_legacy() 

1226 self._ensure_embedding_compat() 

1227 self._warn_stale_doc_prefix() 

1228 

1229 if not self._scalar_ready: 

1230 # A serve-only store never ran ingest, where scalar indexes are 

1231 # built; without them the source/chunk_type prefilters full-scan. 

1232 self.ensure_scalar_indexes(blocking=False) 

1233 # A repair replaces indexes; the handle must see the new versions. 

1234 table.checkout_latest() 

1235 

1236 if query_text: 

1237 hits = self._keyword_arm( 

1238 table, query_text, query_vector, top_k, max_distance, chunk_type 

1239 ) 

1240 if hits is not None: 

1241 return hits 

1242 

1243 try: 

1244 rows = self._vector_arm( 

1245 table, query_vector, top_k * self._config.candidate_multiplier, chunk_type 

1246 ) 

1247 except Exception: 

1248 # A dead scalar index fails the prefilter; the next query re-verifies. 

1249 self._scalar_ready = False 

1250 raise 

1251 log.debug( 

1252 "Vector search: query=%r, candidates=%d, max_distance=%.2f", 

1253 query_text or "vector-only", 

1254 len(rows), 

1255 max_distance, 

1256 ) 

1257 if rows: 

1258 log.debug("Top 5 distances: %s", [r.distance for r in rows[:5]]) 

1259 results = self._filter_and_rerank(rows, query_vector, top_k, max_distance) 

1260 return [ 

1261 r.model_copy( 

1262 update={"score": vector_similarity(r.distance) if r.distance is not None else 0.0} 

1263 ) 

1264 for r in results 

1265 ] 

1266 

1267 def _keyword_arm( 

1268 self, 

1269 table: LanceTable, 

1270 query_text: str, 

1271 query_vector: Vector, 

1272 top_k: int, 

1273 max_distance: float, 

1274 chunk_type: ChunkType | None, 

1275 ) -> list[SearchChunk] | None: 

1276 """Hybrid results, or None when keyword search cannot serve this query and 

1277 the caller must fall back to vector-only. 

1278 

1279 Owns the keyword-search health the endpoint reports. An index that will 

1280 not build and an index that fails mid-query drop every query to vector 

1281 recall alike, so both stamp the flag, and a working index clears it. 

1282 """ 

1283 if not self._fts_ready: 

1284 self.ensure_fts_index(blocking=False) 

1285 table.checkout_latest() 

1286 if self._config.title_search and not self._title_fts_ready: 

1287 self.ensure_title_fts_index(blocking=False) 

1288 if not self._fts_ready: 

1289 # An empty corpus has nothing to index and nothing to degrade. 

1290 self._fts_degraded = self._has_chunks() 

1291 return None 

1292 try: 

1293 hits = self._hybrid_search( 

1294 table, query_text, query_vector, top_k, max_distance, chunk_type 

1295 ) 

1296 except Exception: 

1297 # Falling back changes recall characteristics for the query; 

1298 # a corpus-wide FTS breakage must not present as silence. 

1299 self._fts_degraded = True 

1300 # A dead scalar index fails the same way; the next query re-verifies. 

1301 self._scalar_ready = False 

1302 log.warning("Hybrid search failed, falling back to vector-only", exc_info=True) 

1303 return None 

1304 self._fts_degraded = False 

1305 return hits 

1306 

1307 def _vector_arm( 

1308 self, 

1309 table: LanceTable, 

1310 query_vector: Vector, 

1311 limit: int, 

1312 chunk_type: ChunkType | None, 

1313 ) -> list[SearchChunk]: 

1314 """Vector-arm candidates with the ANN recall recovery applied. 

1315 

1316 When ``chunk_type`` is set, the predicate is pushed into the query so 

1317 the limit applies *after* the type filter; post-filtering would 

1318 silently starve wiki-only queries whose matches live past the window. 

1319 """ 

1320 query = table.search(query_vector).distance_type(_VECTOR_METRIC).limit(limit) 

1321 indices = _index_registry(table) 

1322 if indices is None or _has_vector_index(indices): 

1323 # IVF_PQ is lossy; probe more partitions and refine against full 

1324 # vectors so recall stays close to the exact flat scan. An unreadable 

1325 # registry may hide the index, and the tuning is harmless on a flat scan. 

1326 query = query.nprobes(_ann_nprobes(table.count_rows())) 

1327 query = query.refine_factor(_ANN_REFINE_FACTOR) 

1328 if chunk_type: 

1329 query = query.where(_chunk_type_predicate(chunk_type)) 

1330 return [SearchChunk(**r) for r in query.to_list()] 

1331 

1332 def _fts_arm( 

1333 self, 

1334 table: LanceTable, 

1335 query_text: str, 

1336 limit: int, 

1337 chunk_type: ChunkType | None, 

1338 ) -> list[SearchChunk]: 

1339 """BM25-arm candidates over the chunk text.""" 

1340 return _lexical_rows(table, query_text, limit, chunk_type) 

1341 

1342 def _title_arm( 

1343 self, 

1344 table: LanceTable, 

1345 query_text: str, 

1346 limit: int, 

1347 chunk_type: ChunkType | None, 

1348 ) -> list[SearchChunk]: 

1349 """One BM25 row per document whose title matches, in title-relevance order. 

1350 

1351 Every chunk of a document carries the same title, so all of its chunks 

1352 tie on BM25 and a plain ``limit`` would return an arbitrary tie-ordered 

1353 subset of a single document. Instead this over-fetches, collapses each 

1354 source to one deterministic representative (its first chunk), and returns 

1355 the top *limit* documents ordered by title score -- so "a query naming a 

1356 document by title surfaces its chunks" holds as one stable row per doc. 

1357 

1358 Empty when the store predates the title column or its FTS index (old 

1359 indexes keep working), when the index registry is unreadable, and on 

1360 any query-time failure: the optional 

1361 title arm must never take down the healthy chunk arm, so its failure 

1362 degrades to no-titles, mirroring ``bm25_probe``. 

1363 """ 

1364 indices = _index_registry(table) 

1365 if indices is None or not _has_fts_index(indices, _TITLE_COLUMN): 

1366 return [] 

1367 # Every chunk of one document ties on title BM25, so a fixed window can 

1368 # fill up with a single long document's chunks and starve every other 

1369 # title-matching document. Widen the fetch until enough distinct 

1370 # documents surface, the matches run out, or the ceiling is hit. 

1371 fetch = max(limit * _TITLE_FETCH_FACTOR, _TITLE_MIN_FETCH) 

1372 while True: 

1373 try: 

1374 rows = _lexical_rows(table, query_text, fetch, chunk_type, column=_TITLE_COLUMN) 

1375 except Exception: 

1376 log.debug("Title arm search failed; contributing no title rows", exc_info=True) 

1377 return [] 

1378 best: dict[str, SearchChunk] = {} 

1379 for row in rows: 

1380 seen = best.get(row.source) 

1381 if seen is None or row.chunk_index < seen.chunk_index: 

1382 best[row.source] = row 

1383 if len(best) >= limit or len(rows) < fetch or fetch >= _TITLE_FETCH_CEILING: 

1384 break 

1385 fetch = min(fetch * 4, _TITLE_FETCH_CEILING) 

1386 ordered = sorted(best.values(), key=lambda r: (-(r.bm25_score or 0.0), r.source)) 

1387 return ordered[:limit] 

1388 

1389 def _hybrid_search( 

1390 self, 

1391 table: LanceTable, 

1392 query_text: str, 

1393 query_vector: Vector, 

1394 top_k: int, 

1395 max_distance: float, 

1396 chunk_type: ChunkType | None = None, 

1397 ) -> list[SearchChunk]: 

1398 """Multi-arm retrieval fused by weighted reciprocal rank; the fused ordering is final. 

1399 

1400 A vector arm and a chunk-BM25 arm always run; a title-BM25 arm joins 

1401 when ``cfg.title_search`` is on. Each row's fused score is the 

1402 weight-normalized sum of its arm contributions: the vector arm has 

1403 weight 1.0, the lexical arm ``cfg.lexical_fusion_weight`` (scaled per 

1404 query when ``cfg.adaptive_fusion`` is on), the title arm 

1405 ``cfg.title_search_weight``. So a row a single peaked arm is certain 

1406 about scores that arm's share of the total weight, not a fixed 0.5. 

1407 

1408 Each arm fetches exactly ``top_k`` rows. Deeper pools measurably hurt 

1409 rank fusion by flooding the fused top-k with both-arm mediocrity and 

1410 burying single-arm certainty (lexical identifier hits above all). No 

1411 MMR runs here: lexical passages are often mutually similar, which MMR 

1412 penalizes, trading relevant hits for off-topic neighbors. 

1413 

1414 Title rows carry ``bm25_score``, so a title match counts as lexical 

1415 support for the distance exemption like any other lexical hit. 

1416 """ 

1417 title_rows: list[SearchChunk] = [] 

1418 if self._config.title_search: 

1419 title_rows = self._title_arm(table, query_text, top_k, chunk_type) 

1420 vector_rows = self._vector_arm(table, query_vector, top_k, chunk_type) 

1421 base_lexical_weight = self._config.lexical_fusion_weight 

1422 base_title_weight = self._config.title_search_weight 

1423 lexical_weight = base_lexical_weight 

1424 title_weight = base_title_weight 

1425 if self._config.adaptive_fusion: 

1426 # Quiet the lexical arms per query by vector confidence. The title 

1427 # arm is lexical too, so the same factor scales it. 

1428 scale = adaptive_weight_scale(vector_rows, self._config.adaptive_fusion_margin) 

1429 lexical_weight = base_lexical_weight * scale 

1430 title_weight = base_title_weight * scale 

1431 fused = fuse_arms( 

1432 vector_rows, 

1433 self._fts_arm(table, query_text, top_k, chunk_type), 

1434 title_rows, 

1435 lexical_weight=lexical_weight, 

1436 title_weight=title_weight, 

1437 ) 

1438 fused = _drop_unsupported_far_rows(fused, max_distance) 

1439 return fused[:top_k] 

1440 

1441 def _filter_and_rerank( 

1442 self, 

1443 results: list[SearchChunk], 

1444 query_vector: Vector, 

1445 top_k: int, 

1446 max_distance: float, 

1447 ) -> list[SearchChunk]: 

1448 """Apply the configured distance filter, then MMR-rerank down to top_k.""" 

1449 if max_distance > 0: 

1450 before = len(results) 

1451 if self._config.adaptive_threshold: 

1452 results = self._adaptive_filter(results, top_k, max_distance) 

1453 filter_name = "adaptive" 

1454 else: 

1455 results = self._fixed_filter(results, max_distance) 

1456 filter_name = "fixed" 

1457 log.debug( 

1458 "After %s filter: %d/%d results, threshold=%.2f", 

1459 filter_name, 

1460 len(results), 

1461 before, 

1462 max_distance, 

1463 ) 

1464 if len(results) > top_k: 

1465 results = mmr_rerank(query_vector, results, top_k, self._config.mmr_lambda) 

1466 return results 

1467 

1468 def _adaptive_filter( 

1469 self, results: list[SearchChunk], top_k: int, initial_threshold: float 

1470 ) -> list[SearchChunk]: 

1471 """Widen cosine distance threshold when too few results. 

1472 Inspired by grantflow's (grantflow-ai/grantflow) adaptive retrieval 

1473 pattern which widens thresholds on recursive retry. Step size and 

1474 cap are configurable via ``self._config.adaptive_threshold_step``. 

1475 

1476 Pre-sorts results by distance for a single-pass cutoff search. 

1477 Step size is ``self._config.adaptive_threshold_step`` (default 0.2). 

1478 """ 

1479 cap = max(initial_threshold, _MAX_THRESHOLD) 

1480 step = self._config.adaptive_threshold_step 

1481 

1482 sorted_results = sorted(results, key=_get_distance) 

1483 

1484 threshold = initial_threshold 

1485 for _ in range(_MAX_FILTER_ITERATIONS): 

1486 if threshold > cap: 

1487 break 

1488 cutoff = _count_within_threshold(sorted_results, threshold) 

1489 if cutoff >= top_k: 

1490 return sorted_results[:cutoff] 

1491 threshold += step 

1492 # Final pass at cap 

1493 cutoff = _count_within_threshold(sorted_results, cap) 

1494 return sorted_results[:cutoff] 

1495 

1496 def _fixed_filter(self, results: list[SearchChunk], threshold: float) -> list[SearchChunk]: 

1497 """Simple fixed threshold filter - keep only results within distance threshold.""" 

1498 return [r for r in results if _get_distance(r) <= threshold] 

1499 

1500 def add_entities(self, records: list[dict]) -> int: 

1501 """Append typed entity rows; creates the table on first write. 

1502 

1503 Additive to the store: existing tables and schemas are untouched, and 

1504 stores without this table behave as if nothing was ever extracted. 

1505 """ 

1506 if not records: 

1507 return 0 

1508 from lilbee.retrieval.entities.schema import _entities_schema 

1509 

1510 with self._write_lock(): 

1511 db = self.get_db() 

1512 table = ensure_table(db, ENTITIES_TABLE, _entities_schema()) 

1513 table.add(records) 

1514 return len(records) 

1515 

1516 def entity_schema_state(self) -> EntitySchemaState | None: 

1517 """The persisted entity schema row, or ``None`` when never induced. 

1518 

1519 The schema is machine state induced from the corpus and lives inside 

1520 the index, so it travels with the data. 

1521 """ 

1522 table = self.open_table(ENTITY_SCHEMA_TABLE) 

1523 if table is None: 

1524 return None 

1525 rows = table.search().limit(None).to_list() 

1526 if not rows: 

1527 return None 

1528 # One row by contract; take the newest if a rewrite ever left a stale one. 

1529 row = max(rows, key=lambda r: r["updated_at"]) 

1530 return EntitySchemaState( 

1531 schema_json=str(row["schema_json"]), 

1532 applied=bool(row["applied"]), 

1533 source_count=int(row["source_count"]), 

1534 updated_at=str(row["updated_at"]), 

1535 ) 

1536 

1537 def save_entity_schema(self, schema_json: str, *, applied: bool, source_count: int) -> None: 

1538 """Overwrite the single persisted entity schema row.""" 

1539 with self._write_lock(): 

1540 db = self.get_db() 

1541 table = ensure_table(db, ENTITY_SCHEMA_TABLE, _entity_schema_state_schema()) 

1542 _safe_delete_unlocked(table, ENTITY_SCHEMA_DELETE_ALL_PREDICATE) 

1543 table.add( 

1544 [ 

1545 { 

1546 "schema_json": schema_json, 

1547 "applied": applied, 

1548 "source_count": source_count, 

1549 "updated_at": datetime.now(UTC).isoformat(), 

1550 } 

1551 ] 

1552 ) 

1553 

1554 def mark_entity_schema_applied(self) -> None: 

1555 """Record that a full extraction pass completed under the stored schema.""" 

1556 state = self.entity_schema_state() 

1557 if state is None: 

1558 return 

1559 self.save_entity_schema( 

1560 state["schema_json"], applied=True, source_count=state["source_count"] 

1561 ) 

1562 

1563 def entity_value_counts(self, entity_type: str) -> tuple[int, int]: 

1564 """(mentions, distinct normalized values) for one entity type. 

1565 

1566 Full scan by design: a count is a corpus property. Streaming batches 

1567 keep memory flat at any corpus size. 

1568 """ 

1569 table = self.open_table(ENTITIES_TABLE) 

1570 if table is None: 

1571 return 0, 0 

1572 mentions = 0 

1573 values: set[str] = set() 

1574 arrow = table.to_arrow().select(["type", "normalized_value"]) 

1575 for batch in arrow.to_batches(max_chunksize=_TERM_SCAN_BATCH_ROWS): 

1576 types = batch.column("type").to_pylist() 

1577 vals = batch.column("normalized_value").to_pylist() 

1578 for t_, v in zip(types, vals, strict=True): 

1579 if t_ == entity_type: 

1580 mentions += 1 

1581 values.add(v) 

1582 return mentions, len(values) 

1583 

1584 def entity_association_counts(self, counted: str, grouped_by: str) -> dict[str, int]: 

1585 """Distinct *counted*-type values co-occurring with each *grouped_by* value. 

1586 

1587 Co-occurrence is per chunk: two entities extracted from the same 

1588 ``(source, chunk_index)`` are associated. This is the GROUP BY that 

1589 answers "how many X is each Y associated with". 

1590 """ 

1591 table = self.open_table(ENTITIES_TABLE) 

1592 if table is None: 

1593 return {} 

1594 per_chunk: dict[tuple[str, int], tuple[set[str], set[str]]] = {} 

1595 arrow = table.to_arrow().select(["type", "normalized_value", "source", "chunk_index"]) 

1596 for batch in arrow.to_batches(max_chunksize=_TERM_SCAN_BATCH_ROWS): 

1597 rows = zip( 

1598 batch.column("type").to_pylist(), 

1599 batch.column("normalized_value").to_pylist(), 

1600 batch.column("source").to_pylist(), 

1601 batch.column("chunk_index").to_pylist(), 

1602 strict=True, 

1603 ) 

1604 for t_, v, src, idx in rows: 

1605 if t_ not in (counted, grouped_by): 

1606 continue 

1607 counted_vals, group_vals = per_chunk.setdefault((src, idx), (set(), set())) 

1608 (counted_vals if t_ == counted else group_vals).add(v) 

1609 associations: dict[str, set[str]] = {} 

1610 for counted_vals, group_vals in per_chunk.values(): 

1611 for group_value in group_vals: 

1612 associations.setdefault(group_value, set()).update(counted_vals) 

1613 return {k: len(v) for k, v in sorted(associations.items())} 

1614 

1615 def count_term_mentions(self, term: str) -> tuple[int, int]: 

1616 """(matching chunks, distinct matching sources) for a case-insensitive 

1617 substring scan of the WHOLE chunks table. 

1618 

1619 This is deliberately a full scan, not a top-k search: a count is a 

1620 corpus property, and any retrieval shortcut undercounts it. Streaming 

1621 Arrow batches keeps the working set to one batch of text at a time, 

1622 so cost is linear in corpus size and memory stays flat. 

1623 """ 

1624 table = self.open_table(CHUNKS_TABLE) 

1625 if table is None: 

1626 return 0, 0 

1627 needle = term.lower() 

1628 chunk_hits = 0 

1629 sources: set[str] = set() 

1630 arrow = table.to_arrow().select(["source", "chunk"]) 

1631 for batch in arrow.to_batches(max_chunksize=_TERM_SCAN_BATCH_ROWS): 

1632 texts = batch.column("chunk").to_pylist() 

1633 names = batch.column("source").to_pylist() 

1634 for name, text in zip(names, texts, strict=True): 

1635 if text and needle in text.lower(): 

1636 chunk_hits += 1 

1637 sources.add(name) 

1638 return chunk_hits, len(sources) 

1639 

1640 def count_chunks(self) -> int: 

1641 """Total chunks in the store.""" 

1642 table = self.open_table(CHUNKS_TABLE) 

1643 return table.count_rows() if table is not None else 0 

1644 

1645 def get_chunks_by_source(self, source: str) -> list[SearchChunk]: 

1646 """Return every chunk whose ``source`` equals *source*. 

1647 

1648 The database does the filtering, so only the matching rows are read. 

1649 A query failure raises rather than falling back to a whole-table scan: 

1650 a document's chunks are a bounded read, and the scan that would rescue 

1651 it costs the entire index, vectors included, in memory. 

1652 """ 

1653 table = self.open_table(CHUNKS_TABLE) 

1654 if table is None: 

1655 return [] 

1656 escaped = escape_sql_string(source) 

1657 rows = table.search().where(f"source = '{escaped}'").limit(None).to_list() 

1658 return [SearchChunk(**r) for r in rows] 

1659 

1660 def get_chunks_by_indices(self, source: str, indices: Sequence[int]) -> list[SearchChunk]: 

1661 """Return *source*'s chunks whose ``chunk_index`` is in *indices*. 

1662 

1663 Rows come back in ``chunk_index`` order; indices past either end of 

1664 the document are simply absent from the result. Filtering happens in 

1665 the database for the same reason as :meth:`get_chunks_by_source`: 

1666 neighbor expansion runs once per hit source per query, so a 

1667 whole-table rescue would spike memory on the hottest path there is. 

1668 """ 

1669 if not indices: 

1670 return [] 

1671 table = self.open_table(CHUNKS_TABLE) 

1672 if table is None: 

1673 return [] 

1674 escaped = escape_sql_string(source) 

1675 wanted = ", ".join(str(int(i)) for i in indices) 

1676 predicate = f"source = '{escaped}' AND chunk_index IN ({wanted})" 

1677 rows = table.search().where(predicate).limit(None).to_list() 

1678 return sorted((SearchChunk(**r) for r in rows), key=lambda c: c.chunk_index) 

1679 

1680 def _delete_by_sources_unlocked(self, sources: list[str]) -> None: 

1681 """Delete the sources' chunks, page texts, and chunk-concept rows. 

1682 

1683 Caller must hold ``write_lock()``. One ``IN`` delete per table covers 

1684 every source, so a batched flush pays a constant number of predicate 

1685 deletes instead of one set per document. A delete failure propagates: 

1686 swallowed, it would leave every flushed file silently stale; raised, 

1687 the flush fails and the files replan on the next sync. 

1688 """ 

1689 quoted = ", ".join(f"'{escape_sql_string(source)}'" for source in sources) 

1690 for name, column in _PER_SOURCE_TABLES: 

1691 table = self.open_table(name) 

1692 if table is not None: 

1693 table.delete(f"{column} IN ({quoted})") 

1694 

1695 def _delete_by_source_unlocked(self, source: str) -> None: 

1696 """Delete a single source's chunks, page texts, and chunk-concept rows.""" 

1697 self._delete_by_sources_unlocked([source]) 

1698 

1699 def delete_by_source(self, source: str) -> None: 

1700 """Delete a source's chunks and page texts.""" 

1701 with self._write_lock(): 

1702 self._delete_by_source_unlocked(source) 

1703 self._invalidate_source_cache() 

1704 

1705 def add_page_texts(self, records: list[dict]) -> int: 

1706 """Add per-page text rows (no vectors). Returns count added.""" 

1707 if not records: 

1708 return 0 

1709 with self._write_lock(): 

1710 db = self.get_db() 

1711 table = ensure_table(db, PAGE_TEXTS_TABLE, _page_texts_schema()) 

1712 table.add(records) 

1713 return len(records) 

1714 

1715 def get_page_texts(self, source: str | None = None) -> list[PageTextRecord]: 

1716 """Return per-page text rows, all or for a single *source*.""" 

1717 table = self.open_table(PAGE_TEXTS_TABLE) 

1718 if table is None: 

1719 return [] 

1720 query = table.search() 

1721 if source is not None: 

1722 query = query.where(f"source = '{escape_sql_string(source)}'") 

1723 return cast("list[PageTextRecord]", query.limit(None).to_list()) 

1724 

1725 def page_texts_arrow(self, source: str | None = None) -> pa.Table: 

1726 """Return per-page text rows as an Arrow table in a single scan. 

1727 

1728 The columnar sibling of :meth:`get_page_texts`: the export path keeps the 

1729 whole set in Arrow (no per-row Python objects) from read through file 

1730 write. Empty with the canonical schema when the table or *source* is empty. 

1731 """ 

1732 table = self.open_table(PAGE_TEXTS_TABLE) 

1733 if table is None: 

1734 return _page_texts_schema().empty_table() 

1735 query = table.search().select(["source", "page", "text", "content_type"]) 

1736 if source is not None: 

1737 query = query.where(f"source = '{escape_sql_string(source)}'") 

1738 return query.limit(None).to_arrow() 

1739 

1740 def sources_arrow(self) -> pa.Table: 

1741 """Return each tracked source's extraction metadata as an Arrow table. 

1742 

1743 The columnar sibling of :meth:`get_sources`, keyed by ``source`` so it 

1744 joins straight onto the page-text table. ``get_sources`` builds a dict per 

1745 source, which on a corpus of millions of single-page documents costs more 

1746 than the text being exported. An index written before the metadata columns 

1747 existed gets them as nulls rather than missing, so the join has one shape. 

1748 """ 

1749 import pyarrow as pa 

1750 import pyarrow.compute as pc 

1751 

1752 table = self.open_table(SOURCES_TABLE) 

1753 columns = ["source", *SourceMeta._fields] 

1754 if table is None: 

1755 return pa.schema([pa.field(name, pa.utf8()) for name in columns]).empty_table() 

1756 present = [name for name in SourceMeta._fields if name in table.schema.names] 

1757 arrow = table.search().select(["filename", *present]).limit(None).to_arrow() 

1758 arrow = arrow.rename_columns(["source", *present]) 

1759 for name in SourceMeta._fields: 

1760 if name not in present: 

1761 arrow = arrow.append_column(name, pa.nulls(arrow.num_rows, pa.utf8())) 

1762 arrow = arrow.select(columns) 

1763 if pc.count_distinct(arrow.column("source")).as_py() == arrow.num_rows: 

1764 return arrow 

1765 # One row per source, so a caller joining on it cannot fan out. A doubled 

1766 # row (a source re-merged from a shard) would otherwise multiply every page 

1767 # it owns. Last wins, matching the dict this replaced; single-threaded 

1768 # because that is the only execution mode with an ordered aggregate. 

1769 grouped = arrow.group_by("source", use_threads=False).aggregate( 

1770 [(name, "last") for name in SourceMeta._fields] 

1771 ) 

1772 renamed = grouped.rename_columns( 

1773 [name.removesuffix("_last") for name in grouped.schema.names] 

1774 ) 

1775 return renamed.select(columns) 

1776 

1777 def wiki_chunk_sources(self) -> set[str]: 

1778 """Return the distinct sources of the chunk rows written by the wiki layer.""" 

1779 table = self.open_table(CHUNKS_TABLE) 

1780 if table is None: 

1781 return set() 

1782 rows = ( 

1783 table.search() 

1784 .where(f"chunk_type = '{ChunkType.WIKI}'") 

1785 .select(["source"]) 

1786 .limit(None) 

1787 .to_list() 

1788 ) 

1789 return {row["source"] for row in rows} 

1790 

1791 def wiki_citation_sources(self) -> set[str]: 

1792 """Return the distinct wiki_source values present in the citations table.""" 

1793 table = self.open_table(CITATIONS_TABLE) 

1794 if table is None: 

1795 return set() 

1796 rows = table.search().select(["wiki_source"]).limit(None).to_list() 

1797 return {row["wiki_source"] for row in rows} 

1798 

1799 def get_sources( 

1800 self, 

1801 *, 

1802 search: str | None = None, 

1803 limit: int | None = None, 

1804 offset: int = 0, 

1805 ) -> list[SourceRecord]: 

1806 """Return source records, filtered by *search* and sliced by offset/limit.""" 

1807 table = self.open_table(SOURCES_TABLE) 

1808 if table is None: 

1809 return [] 

1810 query = table.search() 

1811 where = _sources_search_filter(search, include_title="title" in table.schema.names) 

1812 if where is not None: 

1813 query = query.where(where) 

1814 if offset: 

1815 query = query.offset(offset) 

1816 query = query.limit(limit) 

1817 return cast("list[SourceRecord]", query.to_list()) 

1818 

1819 def count_sources(self, *, search: str | None = None) -> int: 

1820 """Count tracked sources matching *search* without materializing rows.""" 

1821 table = self.open_table(SOURCES_TABLE) 

1822 if table is None: 

1823 return 0 

1824 where = _sources_search_filter(search, include_title="title" in table.schema.names) 

1825 count: int = table.count_rows() if where is None else table.count_rows(filter=where) 

1826 return count 

1827 

1828 def _source_row( 

1829 self, 

1830 filename: str, 

1831 file_hash: str, 

1832 chunk_count: int, 

1833 source_type: str, 

1834 stat: SourceStat | None, 

1835 meta: SourceMeta | None = None, 

1836 ) -> dict: 

1837 """Build one ``_sources`` row, defaulting absent stat to the unknown sentinel. 

1838 

1839 Absent extraction metadata persists as NULL, matching rows written 

1840 before the metadata columns existed. 

1841 """ 

1842 meta = meta or SourceMeta() 

1843 return { 

1844 "filename": filename, 

1845 "file_hash": file_hash, 

1846 "ingested_at": datetime.now(UTC).isoformat(), 

1847 "chunk_count": chunk_count, 

1848 "source_type": source_type, 

1849 "size_bytes": stat.size_bytes if stat else SOURCE_STAT_UNKNOWN, 

1850 "mtime_ns": stat.mtime_ns if stat else SOURCE_STAT_UNKNOWN, 

1851 "stat_captured_ns": stat.captured_ns if stat else SOURCE_STAT_UNKNOWN, 

1852 "title": meta.title or None, 

1853 "authors": meta.authors or None, 

1854 "created_at": meta.created_at or None, 

1855 } 

1856 

1857 def _sources_table(self) -> LanceTable: 

1858 """Open/create ``_sources``, adding the stat and metadata columns to older tables.""" 

1859 table = ensure_table(self.get_db(), SOURCES_TABLE, _sources_schema()) 

1860 defaults = {name: f"CAST({SOURCE_STAT_UNKNOWN} AS BIGINT)" for name in _SOURCE_STAT_COLUMNS} 

1861 defaults |= {name: "CAST(NULL AS STRING)" for name in _SOURCE_META_COLUMNS} 

1862 missing = {name: sql for name, sql in defaults.items() if name not in table.schema.names} 

1863 if missing: 

1864 table.add_columns(missing) 

1865 return table 

1866 

1867 def _replace_source_rows_unlocked(self, rows: list[dict]) -> None: 

1868 """Replace source rows: one batched delete plus one batched add. 

1869 

1870 Caller must hold ``write_lock()``. A per-file delete+add pair costs two 

1871 LanceDB version commits, so bulk ingest folds every file in a flush into 

1872 a single pair. 

1873 """ 

1874 table = self._sources_table() 

1875 filenames = ", ".join(f"'{escape_sql_string(r['filename'])}'" for r in rows) 

1876 # Skip the add when the delete failed: adding over a stale row would leave 

1877 # two _sources rows for one filename. The file replans on the next sync. 

1878 if not _safe_delete_unlocked(table, f"filename IN ({filenames})"): 

1879 return 

1880 table.add(rows) 

1881 

1882 def upsert_source( 

1883 self, 

1884 filename: str, 

1885 file_hash: str, 

1886 chunk_count: int, 

1887 source_type: SourceType = SourceType.DOCUMENT, 

1888 stat: SourceStat | None = None, 

1889 meta: SourceMeta | None = None, 

1890 ) -> None: 

1891 """Add or update a source tracking record.""" 

1892 row = self._source_row(filename, file_hash, chunk_count, source_type, stat, meta) 

1893 with self._write_lock(): 

1894 self._replace_source_rows_unlocked([row]) 

1895 self._invalidate_source_cache() 

1896 

1897 def update_source_stats(self, backfills: list[SourceStatBackfill]) -> None: 

1898 """Record size/mtime for already-tracked sources in batched locked writes.""" 

1899 if not backfills: 

1900 return 

1901 for start in range(0, len(backfills), _SOURCE_STAT_BATCH_ROWS): 

1902 rows = [ 

1903 { 

1904 **bf.record, 

1905 "size_bytes": bf.stat.size_bytes, 

1906 "mtime_ns": bf.stat.mtime_ns, 

1907 "stat_captured_ns": bf.stat.captured_ns, 

1908 } 

1909 for bf in backfills[start : start + _SOURCE_STAT_BATCH_ROWS] 

1910 ] 

1911 with self._write_lock(): 

1912 self._replace_source_rows_unlocked(rows) 

1913 self._invalidate_source_cache() 

1914 

1915 def optimize_sources(self) -> None: 

1916 """Compact the sources table; per-flush upserts otherwise accrete tiny versions.""" 

1917 with self._write_lock(): 

1918 table = self.open_table(SOURCES_TABLE) 

1919 if table is None: 

1920 return 

1921 try: 

1922 table.optimize() 

1923 except Exception: 

1924 log.debug("Sources table optimize failed", exc_info=True) 

1925 

1926 def write_chunks_batch(self, items: list[ChunkWrite]) -> int: 

1927 """Write several documents' chunks in one locked transaction. Returns chunks added. 

1928 

1929 One ``write_lock`` acquisition covers the batch's cleanup deletes, page 

1930 texts, chunk add, and source upserts, so a reader never observes a 

1931 half-applied batch. Page texts land after the cleanup and before the 

1932 source rows, so a page-text failure leaves the rows stale and the files 

1933 replan next sync; a document with no chunks still persists its page 

1934 texts and source row. The embedding-identity gate and per-vector 

1935 dimension check mirror ``add_chunks``; a dimension mismatch raises and 

1936 the whole batch is rejected. 

1937 """ 

1938 if not items: 

1939 return 0 

1940 with self._write_lock(timeout=BATCH_LOCK_TIMEOUT): 

1941 embedding_model = self._config.embedding_model 

1942 embedding_dim = self._config.embedding_dim 

1943 self._ensure_embedding_compat() 

1944 self._fts_ready = False 

1945 self._scalar_ready = False 

1946 all_records = [rec for it in items for rec in it.records] 

1947 _check_vector_dims(all_records, embedding_dim) 

1948 db = self.get_db() 

1949 self._cleanup_batch_unlocked(items) 

1950 self._add_page_texts_unlocked(db, items) 

1951 self._add_chunk_records_unlocked(all_records, embedding_model, embedding_dim) 

1952 self._replace_source_rows_unlocked(self._batch_source_rows(items)) 

1953 self._invalidate_source_cache() 

1954 return len(all_records) 

1955 

1956 def _cleanup_batch_unlocked(self, items: list[ChunkWrite]) -> None: 

1957 """One ``IN`` delete per table for the flagged documents. Caller holds ``write_lock()``.""" 

1958 cleanup_sources = [it.source for it in items if it.needs_cleanup] 

1959 if cleanup_sources: 

1960 self._delete_by_sources_unlocked(cleanup_sources) 

1961 

1962 def _add_page_texts_unlocked(self, db: LanceDBConnection, items: list[ChunkWrite]) -> None: 

1963 """Add the batch's page-text rows. Caller holds ``write_lock()``.""" 

1964 page_rows = [row for it in items for row in (it.page_texts or [])] 

1965 if page_rows: 

1966 ensure_table(db, PAGE_TEXTS_TABLE, _page_texts_schema()).add(page_rows) 

1967 

1968 def _add_chunk_records_unlocked( 

1969 self, 

1970 all_records: list[dict], 

1971 embedding_model: str, 

1972 embedding_dim: int, 

1973 ) -> None: 

1974 """Add the batch's chunk rows, writing meta on first use. Caller holds ``write_lock()``.""" 

1975 if not all_records: 

1976 return 

1977 self._chunks_table().add(all_records) 

1978 if self.get_meta() is None: 

1979 self._write_meta_unlocked(embedding_model=embedding_model, embedding_dim=embedding_dim) 

1980 

1981 def _batch_source_rows(self, items: list[ChunkWrite]) -> list[dict]: 

1982 """One ``_sources`` row per batched document.""" 

1983 return [ 

1984 self._source_row( 

1985 it.source, it.file_hash, len(it.records), it.source_type, it.stat, it.meta 

1986 ) 

1987 for it in items 

1988 ] 

1989 

1990 def _delete_source_unlocked(self, filename: str) -> None: 

1991 """Remove the *filename* source record. Caller must hold ``write_lock()``.""" 

1992 table = self.open_table(SOURCES_TABLE) 

1993 if table is not None: 

1994 _safe_delete_unlocked(table, f"filename = '{escape_sql_string(filename)}'") 

1995 

1996 def delete_source(self, filename: str) -> None: 

1997 """Remove a source file tracking record.""" 

1998 with self._write_lock(): 

1999 self._delete_source_unlocked(filename) 

2000 self._invalidate_source_cache() 

2001 

2002 def _remove_many_unlocked(self, names: list[str]) -> None: 

2003 """Delete the documents' chunks and source records together. 

2004 

2005 All deletes run under the caller's single ``write_lock()`` so no 

2006 reader can observe chunks whose source record is already gone; one 

2007 ``IN`` delete per table covers the whole set. 

2008 """ 

2009 self._delete_by_sources_unlocked(names) 

2010 quoted = ", ".join(f"'{escape_sql_string(name)}'" for name in names) 

2011 table = self.open_table(SOURCES_TABLE) 

2012 if table is not None: 

2013 _safe_delete_unlocked(table, f"filename IN ({quoted})") 

2014 

2015 def relocate_sources(self, moves: list[tuple[str, str, SourceStat | None]]) -> None: 

2016 """Re-key moved sources from old filename to new, preserving their chunks. 

2017 

2018 A source whose file moved (same content hash, new path) keeps its chunks 

2019 and embeddings; only its filename key and disk stat change. Each per-source 

2020 table's source column, the citation source_filename, and the sources row are 

2021 updated in place under one write lock, so a move costs no re-extraction or 

2022 re-embedding. ``moves`` is ``(old_name, new_name, new_stat)`` tuples. 

2023 

2024 Each table is opened once; the re-key is then a targeted per-move update. 

2025 A single-statement batch would need a ``CASE`` expression, which LanceDB's 

2026 update SQL does not support, and a delete+re-add across the vector tables is 

2027 not worth its risk for what is a rare mass relabel. 

2028 """ 

2029 if not moves: 

2030 return 

2031 from lilbee.data.title import derive_title # circular at module scope 

2032 

2033 with self._write_lock(): 

2034 tables = [(self.open_table(name), column) for name, column in _RELOCATABLE_TABLES] 

2035 sources = self.open_table(SOURCES_TABLE) 

2036 for old, new, stat in moves: 

2037 where_old = f"= '{escape_sql_string(old)}'" 

2038 new_title = self._relocated_title(sources, old, new, derive_title) 

2039 for table, column in tables: 

2040 if table is None: 

2041 continue 

2042 values: dict[str, object] = {column: new} 

2043 # Stem titles track the filename; re-derive them on the same 

2044 # handle and statement as the re-key. 

2045 if new_title is not _KEEP_TITLE and _TITLE_COLUMN in table.schema.names: 

2046 values[_TITLE_COLUMN] = new_title 

2047 table.update(where=f"{column} {where_old}", values=values) 

2048 if sources is not None: 

2049 row_values: dict[str, object] = {"filename": new} 

2050 if new_title is not _KEEP_TITLE: 

2051 row_values["title"] = new_title 

2052 if stat is not None: 

2053 row_values["size_bytes"] = stat.size_bytes 

2054 row_values["mtime_ns"] = stat.mtime_ns 

2055 row_values["stat_captured_ns"] = stat.captured_ns 

2056 sources.update(where=f"filename {where_old}", values=row_values) 

2057 self._invalidate_source_cache() 

2058 

2059 def _relocated_title( 

2060 self, 

2061 sources: LanceTable | None, 

2062 old: str, 

2063 new: str, 

2064 derive: Callable[[str], str], 

2065 ) -> str | None: 

2066 """New title for a moved source, or ``_KEEP_TITLE`` when it must not change. 

2067 

2068 Extraction-derived titles survive a move (the content is unchanged); 

2069 a stem-derived title tracks the filename it was derived from, so it is 

2070 re-derived from the new name instead of matching the old one forever. 

2071 """ 

2072 if sources is None: 

2073 return _KEEP_TITLE 

2074 try: 

2075 rows = ( 

2076 sources.search() 

2077 .where(f"filename = '{escape_sql_string(old)}'") 

2078 .select(["title"]) 

2079 .limit(1) 

2080 .to_list() 

2081 ) 

2082 except Exception: 

2083 return _KEEP_TITLE 

2084 if not rows or "title" not in rows[0]: 

2085 return _KEEP_TITLE 

2086 stored = rows[0]["title"] or "" 

2087 if stored != (derive(old) or ""): 

2088 return _KEEP_TITLE 

2089 return derive(new) or None 

2090 

2091 def member_sources(self, name: str) -> list[str]: 

2092 """Sources ingested out of the archive *name*: every filename under ``name/``.""" 

2093 prefix = f"{name}/" 

2094 return [s["filename"] for s in self.get_sources() if s["filename"].startswith(prefix)] 

2095 

2096 def remove_documents(self, names: list[str]) -> RemoveResult: 

2097 """Remove documents from the knowledge base by source name. 

2098 

2099 Looks up known sources and deletes their chunks and source records. Never 

2100 touches files on disk: source bytes are the user's, and a linked-in corpus 

2101 must never be deleted. Durable, file-aware removal (skip-markers, unlinking 

2102 a top-level link) lives in :func:`lilbee.app.ingest.remove_documents_durably`. 

2103 

2104 Returns a RemoveResult with removed and not_found lists. 

2105 """ 

2106 known = {s["filename"] for s in self.get_sources()} 

2107 targets = [*names, *(m for name in names for m in self.member_sources(name))] 

2108 removed = [name for name in targets if name in known] 

2109 not_found = [name for name in names if name not in known] 

2110 

2111 if removed: 

2112 # One lock acquisition and one IN-delete per table for the whole 

2113 # set, mirroring the batched flush path, instead of a LanceDB 

2114 # version commit per document. 

2115 with self._write_lock(): 

2116 self._remove_many_unlocked(removed) 

2117 self._invalidate_source_cache() 

2118 

2119 return RemoveResult(removed=removed, not_found=not_found) 

2120 

2121 def clear_table(self, name: str, predicate: str) -> bool: 

2122 """Delete rows matching *predicate* from *name*. Acquires write lock. 

2123 

2124 Returns whether the delete succeeded, so a caller recording the 

2125 outcome (prune, the legacy migration) does not report success over a 

2126 swallowed failure. 

2127 """ 

2128 with self._write_lock(): 

2129 table = self.open_table(name) 

2130 if table is None: 

2131 return True 

2132 return _safe_delete_unlocked(table, predicate) 

2133 

2134 def clear_and_add(self, name: str, schema: pa.Schema, rows: list[dict], predicate: str) -> None: 

2135 """Replace the rows matching *predicate* with *rows* in one locked write. 

2136 

2137 Delete and add run under a single write lock, so a reader never observes 

2138 the table emptied mid-rebuild. A delete failure propagates, following the 

2139 same rule as :meth:`_delete_by_sources_unlocked`: adding over rows whose 

2140 predecessors are still there duplicates them, and reporting success would 

2141 leave the caller acting on state it thinks it replaced. 

2142 """ 

2143 with self._write_lock(): 

2144 db = self.get_db() 

2145 table = ensure_table(db, name, schema) 

2146 table.delete(predicate) 

2147 if rows: 

2148 table.add(rows) 

2149 

2150 def add_citations(self, records: list[CitationRecord]) -> int: 

2151 """Add citation records to the store. Returns count added.""" 

2152 if not records: 

2153 return 0 

2154 with self._write_lock(): 

2155 db = self.get_db() 

2156 table = ensure_table(db, CITATIONS_TABLE, _citations_schema()) 

2157 table.add(records) 

2158 return len(records) 

2159 

2160 def get_citations_for_wiki(self, wiki_source: str) -> list[CitationRecord]: 

2161 """Get all citations for a wiki page.""" 

2162 table = self.open_table(CITATIONS_TABLE) 

2163 if table is None: 

2164 return [] 

2165 escaped = escape_sql_string(wiki_source) 

2166 rows = table.search().where(f"wiki_source = '{escaped}'").to_list() 

2167 return cast("list[CitationRecord]", rows) 

2168 

2169 def get_citations_for_source(self, source_filename: str) -> list[CitationRecord]: 

2170 """Get all citations that reference a source document (reverse lookup).""" 

2171 table = self.open_table(CITATIONS_TABLE) 

2172 if table is None: 

2173 return [] 

2174 escaped = escape_sql_string(source_filename) 

2175 rows = table.search().where(f"source_filename = '{escaped}'").to_list() 

2176 return cast("list[CitationRecord]", rows) 

2177 

2178 def delete_citations_for_wiki(self, wiki_source: str) -> bool: 

2179 """Delete all citations for a wiki page. Returns whether the delete succeeded.""" 

2180 return self.clear_table(CITATIONS_TABLE, _citations_for_wiki_predicate(wiki_source)) 

2181 

2182 def delete_all_wiki_rows(self) -> bool: 

2183 """Delete every wiki chunk row, every citation, and every mention. 

2184 Returns whether all deletes succeeded, so a caller cannot report a wipe 

2185 over a swallowed failure. Only the wiki layer writes citations and 

2186 mentions, so wiping it empties those tables outright. All deletes run 

2187 before the results combine. 

2188 """ 

2189 chunks_cleared = self.clear_table(CHUNKS_TABLE, f"chunk_type = '{ChunkType.WIKI}'") 

2190 citations_cleared = self.clear_table(CITATIONS_TABLE, "1 = 1") 

2191 mentions_cleared = self.clear_wiki_mentions() 

2192 return chunks_cleared and citations_cleared and mentions_cleared 

2193 

2194 def replace_citations_for_wiki(self, wiki_source: str, records: list[CitationRecord]) -> None: 

2195 """Swap a wiki page's citations for *records* under one write lock. 

2196 

2197 The lock keeps another writer from landing between the delete and the 

2198 add; the two remain separate commits, so a crash in between leaves the 

2199 rows absent until the page is regenerated or accepted again. 

2200 """ 

2201 self.clear_and_add( 

2202 CITATIONS_TABLE, 

2203 _citations_schema(), 

2204 [dict(rec) for rec in records], 

2205 _citations_for_wiki_predicate(wiki_source), 

2206 ) 

2207 

2208 def replace_wiki_mentions_for_source(self, source: str, rows: list[dict]) -> None: 

2209 """Swap one source's wiki mention rows for *rows* under one write lock. 

2210 

2211 A source contributes its whole mention set at once, so an incremental 

2212 wiki refresh replaces exactly that source's rows and leaves every other 

2213 source's evidence in place for the corpus-wide aggregate. 

2214 """ 

2215 escaped = escape_sql_string(source) 

2216 self.clear_and_add( 

2217 WIKI_MENTIONS_TABLE, 

2218 _wiki_mentions_schema(), 

2219 rows, 

2220 f"source = '{escaped}'", 

2221 ) 

2222 

2223 def wiki_mention_rows(self, slugs: Iterable[str] | None = None) -> list[dict]: 

2224 """Every wiki mention row, or only those for *slugs*. 

2225 

2226 The wiki aggregates these across sources to rebuild the stub index. A 

2227 full refresh reads them all; an incremental one reads only the slugs it 

2228 touched, to recompute their corpus-wide totals without a full scan. 

2229 """ 

2230 table = self.open_table(WIKI_MENTIONS_TABLE) 

2231 if table is None: 

2232 return [] 

2233 query = table.search() 

2234 if slugs is not None: 

2235 wanted = list(slugs) 

2236 if not wanted: 

2237 return [] 

2238 joined = ", ".join(f"'{escape_sql_string(s)}'" for s in wanted) 

2239 query = query.where(f"slug IN ({joined})") 

2240 rows: list[dict] = query.limit(None).to_list() 

2241 return rows 

2242 

2243 def clear_wiki_mentions(self) -> bool: 

2244 """Drop every wiki mention row (a full rebuild starts from nothing).""" 

2245 return self.clear_table(WIKI_MENTIONS_TABLE, "1 = 1") 

2246 

2247 def has_wiki_mentions(self) -> bool: 

2248 """Whether any wiki mention row exists. 

2249 

2250 A cold store or one migrated from the file-only index has none, which 

2251 forces a refresh to rebuild in full and seed the table before an 

2252 incremental pass can aggregate over it. 

2253 """ 

2254 table = self.open_table(WIKI_MENTIONS_TABLE) 

2255 return table is not None and table.count_rows() > 0 

2256 

2257 def _memories_schema(self) -> pa.Schema: 

2258 return pa.schema( 

2259 [ 

2260 pa.field("id", pa.utf8()), 

2261 pa.field("owner", pa.utf8()), 

2262 pa.field("shared", pa.bool_()), 

2263 pa.field("kind", pa.utf8()), 

2264 pa.field("source", pa.utf8()), 

2265 pa.field("text", pa.utf8()), 

2266 pa.field("vector", pa.list_(pa.float32(), self._config.embedding_dim)), 

2267 pa.field("created_at", pa.utf8()), 

2268 pa.field("updated_at", pa.utf8()), 

2269 ] 

2270 ) 

2271 

2272 def _duplicate_memory_id_unlocked(self, table: LanceTable, record: MemoryRow) -> str | None: 

2273 """Return the id of a near-duplicate same-owner, same-kind memory, if any.""" 

2274 if table.count_rows() == 0: 

2275 return None 

2276 predicate = ( 

2277 f"owner = '{escape_sql_string(record.owner)}' " 

2278 f"AND kind = '{escape_sql_string(record.kind)}'" 

2279 ) 

2280 rows = ( 

2281 table.search(record.vector) 

2282 .distance_type(_VECTOR_METRIC) 

2283 .where(predicate) 

2284 .limit(1) 

2285 .to_list() 

2286 ) 

2287 if rows and rows[0].get("_distance", 1.0) <= self._config.memory_dedup_distance: 

2288 return str(rows[0]["id"]) 

2289 return None 

2290 

2291 def _evict_overflow_unlocked(self, table: LanceTable, owner: str) -> None: 

2292 """Delete oldest memories for *owner* so an incoming insert stays within the cap.""" 

2293 cap = self._config.memory_max_per_owner 

2294 predicate = f"owner = '{escape_sql_string(owner)}'" 

2295 rows = table.search().where(predicate).limit(None).to_list() 

2296 if len(rows) < cap: 

2297 return 

2298 rows.sort(key=lambda r: r.get("created_at", "")) 

2299 for row in rows[: len(rows) - (cap - 1)]: 

2300 _safe_delete_unlocked(table, f"id = '{escape_sql_string(str(row['id']))}'") 

2301 

2302 def add_memory(self, record: MemoryRow) -> str: 

2303 """Insert *record*, or update the nearest same-owner duplicate in place. 

2304 

2305 Returns the stored id. Raises ``EmbeddingModelMismatchError`` when the store 

2306 was built under a different embedding model, and ``ValueError`` on a vector 

2307 dimension mismatch. 

2308 """ 

2309 if len(record.vector) != self._config.embedding_dim: 

2310 raise ValueError( 

2311 f"Memory vector dimension mismatch: expected " 

2312 f"{self._config.embedding_dim}, got {len(record.vector)}" 

2313 ) 

2314 with self._write_lock(): 

2315 embedding_model = self._config.embedding_model 

2316 embedding_dim = self._config.embedding_dim 

2317 self._ensure_embedding_compat() 

2318 db = self.get_db() 

2319 table = ensure_table(db, MEMORIES_TABLE, self._memories_schema()) 

2320 duplicate_id = self._duplicate_memory_id_unlocked(table, record) 

2321 if duplicate_id is not None and _safe_delete_unlocked( 

2322 table, f"id = '{escape_sql_string(duplicate_id)}'" 

2323 ): 

2324 # Only reuse the id once the old row is actually gone; a swallowed 

2325 # delete failure would otherwise leave two rows with the same id. 

2326 record.id = duplicate_id 

2327 self._evict_overflow_unlocked(table, record.owner) 

2328 table.add([record.model_dump(mode="json")]) 

2329 if self.get_meta() is None: 

2330 self._write_meta_unlocked( 

2331 embedding_model=embedding_model, embedding_dim=embedding_dim 

2332 ) 

2333 return record.id 

2334 

2335 def get_memories( 

2336 self, 

2337 *, 

2338 owner_predicate: str, 

2339 kind: MemoryKind | None = None, 

2340 ) -> list[MemoryRow]: 

2341 """Return memories matching *owner_predicate* and optional *kind*, newest first.""" 

2342 table = self.open_table(MEMORIES_TABLE) 

2343 if table is None: 

2344 return [] 

2345 clauses = [f"({owner_predicate})"] 

2346 if kind is not None: 

2347 clauses.append(f"kind = '{escape_sql_string(kind)}'") 

2348 rows = table.search().where(" AND ".join(clauses)).limit(None).to_list() 

2349 memories = [MemoryRow(**r) for r in rows] 

2350 memories.sort(key=lambda m: m.created_at, reverse=True) 

2351 return memories 

2352 

2353 def search_memories( 

2354 self, 

2355 query_vector: Vector, 

2356 *, 

2357 owner_predicate: str, 

2358 top_k: int, 

2359 max_distance: float, 

2360 ) -> list[MemoryRow]: 

2361 """Vector-recall FACT memories within *max_distance*, best first.""" 

2362 table = self.open_table(MEMORIES_TABLE) 

2363 if table is None or top_k <= 0: 

2364 return [] 

2365 self._ensure_embedding_compat() 

2366 predicate = f"({owner_predicate}) AND kind = '{MemoryKind.FACT}'" 

2367 rows = ( 

2368 table.search(query_vector) 

2369 .distance_type(_VECTOR_METRIC) 

2370 .where(predicate) 

2371 .limit(top_k) 

2372 .to_list() 

2373 ) 

2374 return [MemoryRow(**r) for r in rows if r.get("_distance", 1.0) <= max_distance] 

2375 

2376 def update_memory(self, memory_id: str, *, shared: bool, owner: str) -> bool: 

2377 """Set the *shared* flag on *owner*'s memory. Returns True when found and owned. 

2378 

2379 The ``owner`` predicate scopes the mutation to the caller's namespace so an 

2380 agent cannot flip another owner's (or the human's) memory. 

2381 """ 

2382 with self._write_lock(): 

2383 table = self.open_table(MEMORIES_TABLE) 

2384 if table is None: 

2385 return False 

2386 predicate = self._owned_memory_predicate(memory_id, owner) 

2387 rows = table.search().where(predicate).limit(1).to_list() 

2388 if not rows: 

2389 return False 

2390 record = MemoryRow(**rows[0]) 

2391 record.shared = shared 

2392 record.updated_at = datetime.now(UTC).isoformat() 

2393 # If the delete fails, do not add the modified copy: that would leave 

2394 # two rows for one id. Report not-updated instead. 

2395 if not _safe_delete_unlocked(table, predicate): 

2396 return False 

2397 table.add([record.model_dump(mode="json")]) 

2398 return True 

2399 

2400 def delete_memory(self, memory_id: str, *, owner: str) -> bool: 

2401 """Delete *owner*'s memory by id. Returns True when a matching row was deleted. 

2402 

2403 The ``owner`` predicate scopes the delete to the caller's namespace so an 

2404 agent cannot destroy another owner's (or the human's) memory. 

2405 """ 

2406 with self._write_lock(): 

2407 table = self.open_table(MEMORIES_TABLE) 

2408 if table is None: 

2409 return False 

2410 predicate = self._owned_memory_predicate(memory_id, owner) 

2411 if not table.search().where(predicate).limit(1).to_list(): 

2412 return False 

2413 # Report the real outcome: a swallowed delete failure must not be 

2414 # reported as a successful forget. 

2415 return _safe_delete_unlocked(table, predicate) 

2416 

2417 @staticmethod 

2418 def _owned_memory_predicate(memory_id: str, owner: str) -> str: 

2419 """SQL predicate matching a single memory id within *owner*'s namespace.""" 

2420 return f"id = '{escape_sql_string(memory_id)}' AND owner = '{escape_sql_string(owner)}'" 

2421 

2422 def rebuild_memory_embeddings(self, embed: Callable[[list[str]], list[Vector]]) -> int: 

2423 """Re-embed every memory under the current model, recreating the table. 

2424 

2425 The vector column dimension is immutable, so a different-dim model needs a 

2426 fresh table; recreating unconditionally also covers the same-dim case. Memory 

2427 text is human-authored and re-embeddable, so no data is lost. Returns the count. 

2428 

2429 The snapshot, embed, and table rebuild all run under the write lock so a 

2430 concurrent ``add_memory`` cannot commit into the read-then-drop window and 

2431 be erased; it either lands before the snapshot or blocks until the rebuild 

2432 finishes. 

2433 """ 

2434 with self._write_lock(): 

2435 table = self.open_table(MEMORIES_TABLE) 

2436 if table is None: 

2437 return 0 

2438 rows = table.search().limit(None).to_list() 

2439 if not rows: 

2440 return 0 

2441 memories = [MemoryRow(**r) for r in rows] 

2442 vectors = embed([m.text for m in memories]) 

2443 for memory, vector in zip(memories, vectors, strict=True): 

2444 # MemoryRow serializes to JSON on write, which has no ndarray encoding. 

2445 memory.vector = vector.tolist() 

2446 db = self.get_db() 

2447 db.drop_table(MEMORIES_TABLE) 

2448 new_table = ensure_table(db, MEMORIES_TABLE, self._memories_schema()) 

2449 new_table.add([m.model_dump(mode="json") for m in memories]) 

2450 return len(memories) 

2451 

2452 def close(self) -> None: 

2453 """Release the database connection and reset state.""" 

2454 self._db = None 

2455 self._fts_ready = False 

2456 self._title_fts_ready = False 

2457 self._scalar_ready = False 

2458 

2459 def drop_all(self) -> None: 

2460 """Drop every table except ``_memories`` -- used by rebuild. 

2461 

2462 Memory is user-authored data with no on-disk source, not derived from 

2463 documents, so a rebuild preserves it. Only a factory reset (which deletes 

2464 the data directory) clears it. 

2465 """ 

2466 with self._write_lock(): 

2467 self._fts_ready = False 

2468 self._title_fts_ready = False 

2469 self._scalar_ready = False 

2470 db = self.get_db() 

2471 for name in table_names(db): 

2472 if name == MEMORIES_TABLE: 

2473 continue 

2474 db.drop_table(name) 

2475 self._invalidate_source_cache()