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
« 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."""
3from __future__ import annotations
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
14import pyarrow as pa
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
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)
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
92log = logging.getLogger(__name__)
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
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.
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 ]
122_MAX_THRESHOLD = 1.0
123_MAX_FILTER_ITERATIONS = 20 # safety cap to prevent runaway loops
126def _is_fts_position_overflow(exc: Exception) -> bool:
127 """True when *exc* is LanceDB's positional-FTS list-encoding overflow.
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
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*.
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
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()]
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
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")
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")
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"
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)
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)
217# Sentinel: relocation must leave the stored title untouched (extraction-derived).
218_KEEP_TITLE = "\x00keep"
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
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
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
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))
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 )
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)}'"
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")
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)
272class Store:
273 """LanceDB vector store: wraps all DB operations with config-driven defaults."""
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
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
297 return FTS(with_position=False, language=self._config.fts_language)
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)
303 def _write_lock(self, timeout: float = LOCK_TIMEOUT) -> AbstractContextManager[None]:
304 """Acquire the write lock keyed on *this* store's data directory.
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)
312 def _invalidate_source_cache(self) -> None:
313 """Drop the cached {filename: ingested_at} map."""
314 self._source_ingested_cache = None
316 def source_ingested_at_map(self) -> dict[str, str]:
317 """Return {filename: ingested_at} for every source, cached until mutation.
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
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 )
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
355 def _backfill_stem_titles_unlocked(self, table: LanceTable) -> None:
356 """Backfill filename-stem titles for pre-upgrade rows. Caller holds ``write_lock()``.
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
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 )
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 )
406 def _write_meta_unlocked(self, *, embedding_model: str, embedding_dim: int) -> None:
407 """Overwrite the single ``_meta`` row with the supplied identity.
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
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
435 def has_chunks(self) -> bool:
436 """Public predicate: True iff the store currently holds at least one chunk."""
437 return self._has_chunks()
439 def initialize_meta_if_legacy(self) -> bool:
440 """Pin a legacy store's identity to the current cfg if not already set.
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
469 def index_mismatch(self) -> EmbeddingModelMismatchError | None:
470 """The drift between the persisted embedding identity and cfg, or None.
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 )
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
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
512 def _warn_stale_doc_prefix(self) -> None:
513 """Warn once when the embedding family's document prefix postdates this store.
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 )
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 )
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
599 def assert_embedding_compatible(self) -> None:
600 """Run the full embedding-identity gate (legacy init, canonicalize, check).
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()
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 )
620 def canonicalize_meta_if_legacy(self) -> bool:
621 """Rewrite a legacy bare-repo ``_meta`` row to the canonical full ref.
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
648 def get_db(self) -> LanceDBConnection:
649 if self._db is None:
650 from lancedb.db import LanceDBConnection
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
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
667 def ensure_fts_index(self, *, blocking: bool = True) -> None:
668 """Create the chunks FTS index, or run ``optimize()`` once it exists.
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.
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")
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)
737 def _ensure_title_fts_unlocked(self, table: LanceTable) -> None:
738 """Create the title FTS index when the column exists. Caller holds ``write_lock()``.
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 )
764 def ensure_title_fts_index(self, *, blocking: bool = True) -> None:
765 """Build the title FTS index for a title_search toggle after startup.
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")
786 def _rebuild_fts(self, table: LanceTable, reason: str) -> None:
787 """Replace the FTS indexes with fresh positionless ones. Caller holds the lock.
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 )
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.
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.
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
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 )
854 def ensure_scalar_indexes(self, *, blocking: bool = True) -> None:
855 """Build scalar indexes on the columns lilbee filters by.
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")
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*.
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 )
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
943 table.create_index("vector", config=IvfPq(distance_type=_VECTOR_METRIC), replace=True)
945 def ensure_vector_index(self, *, force: bool = False) -> bool:
946 """Build or refresh the ANN vector index when the corpus is large enough.
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)
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
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
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)
1026 def add_chunks(self, records: list[dict]) -> int:
1027 """Add chunk records to the store. Returns count added.
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.
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)
1040 def replace_chunks(self, records: list[dict], predicate: str) -> int:
1041 """Replace the chunk rows matching *predicate* with *records* under one write lock.
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)
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)
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.
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)
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.
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
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 )
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.
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.
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.
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
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
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.
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 []
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.
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.
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()
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()
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
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 ]
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.
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
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.
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()]
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)
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.
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.
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]
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.
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.
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.
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]
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
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``.
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
1482 sorted_results = sorted(results, key=_get_distance)
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]
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]
1500 def add_entities(self, records: list[dict]) -> int:
1501 """Append typed entity rows; creates the table on first write.
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
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)
1516 def entity_schema_state(self) -> EntitySchemaState | None:
1517 """The persisted entity schema row, or ``None`` when never induced.
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 )
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 )
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 )
1563 def entity_value_counts(self, entity_type: str) -> tuple[int, int]:
1564 """(mentions, distinct normalized values) for one entity type.
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)
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.
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())}
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.
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)
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
1645 def get_chunks_by_source(self, source: str) -> list[SearchChunk]:
1646 """Return every chunk whose ``source`` equals *source*.
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]
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*.
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)
1680 def _delete_by_sources_unlocked(self, sources: list[str]) -> None:
1681 """Delete the sources' chunks, page texts, and chunk-concept rows.
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})")
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])
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()
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)
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())
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.
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()
1740 def sources_arrow(self) -> pa.Table:
1741 """Return each tracked source's extraction metadata as an Arrow table.
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
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)
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}
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}
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())
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
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.
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 }
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
1867 def _replace_source_rows_unlocked(self, rows: list[dict]) -> None:
1868 """Replace source rows: one batched delete plus one batched add.
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)
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()
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()
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)
1926 def write_chunks_batch(self, items: list[ChunkWrite]) -> int:
1927 """Write several documents' chunks in one locked transaction. Returns chunks added.
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)
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)
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)
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)
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 ]
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)}'")
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()
2002 def _remove_many_unlocked(self, names: list[str]) -> None:
2003 """Delete the documents' chunks and source records together.
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})")
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.
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.
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
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()
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.
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
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)]
2096 def remove_documents(self, names: list[str]) -> RemoveResult:
2097 """Remove documents from the knowledge base by source name.
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`.
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]
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()
2119 return RemoveResult(removed=removed, not_found=not_found)
2121 def clear_table(self, name: str, predicate: str) -> bool:
2122 """Delete rows matching *predicate* from *name*. Acquires write lock.
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)
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.
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)
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)
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)
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)
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))
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
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.
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 )
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.
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 )
2223 def wiki_mention_rows(self, slugs: Iterable[str] | None = None) -> list[dict]:
2224 """Every wiki mention row, or only those for *slugs*.
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
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")
2247 def has_wiki_mentions(self) -> bool:
2248 """Whether any wiki mention row exists.
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
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 )
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
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']))}'")
2302 def add_memory(self, record: MemoryRow) -> str:
2303 """Insert *record*, or update the nearest same-owner duplicate in place.
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
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
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]
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.
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
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.
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)
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)}'"
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.
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.
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)
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
2459 def drop_all(self) -> None:
2460 """Drop every table except ``_memories`` -- used by rebuild.
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()