use std::sync::atomic::Ordering;
use anyhow::{Context, Result};
use lancedb::index::scalar::FtsIndexBuilder;
use lancedb::index::Index;
use lancedb::table::OptimizeAction;
use lancedb::Table;
use super::{is_commit_conflict, VectorStore};
impl VectorStore {
pub async fn ensure_fts_index(&self) -> Result<()> {
if self.fts_indexed.load(Ordering::Relaxed) {
return Ok(());
}
let table = match &self.table {
Some(t) => t,
None => return Ok(()),
};
if Self::has_content_fts_index(table).await {
self.fts_indexed.store(true, Ordering::Relaxed);
return Ok(());
}
crate::operational_metrics::record_fts_build_attempt();
self.fts_build_attempts.fetch_add(1, Ordering::Relaxed);
let result = table
.create_index(&["content"], Index::FTS(FtsIndexBuilder::default()))
.replace(false)
.execute()
.await;
match result {
Ok(()) => {}
Err(_) => {
}
}
self.fts_indexed.store(true, Ordering::Relaxed);
Ok(())
}
#[cfg(test)]
pub(crate) fn fts_build_attempts(&self) -> u64 {
self.fts_build_attempts.load(Ordering::Relaxed)
}
pub(super) async fn has_content_fts_index(table: &Table) -> bool {
table
.list_indices()
.await
.map(|indices| {
indices.iter().any(|index| {
index.index_type == lancedb::index::IndexType::FTS
&& index.columns.iter().any(|c| c == "content")
})
})
.unwrap_or(false)
}
pub async fn fts_coverage(&self) -> Option<(usize, usize)> {
let table = self.table.as_ref()?;
let indices = table.list_indices().await.ok()?;
let index = indices.iter().find(|index| {
index.index_type == lancedb::index::IndexType::FTS
&& index.columns.iter().any(|c| c == "content")
})?;
let stats = table.index_stats(&index.name).await.ok()??;
Some((stats.num_indexed_rows, stats.num_unindexed_rows))
}
pub async fn rebuild_fts_index(&self) -> Result<()> {
let table = match &self.table {
Some(t) => t,
None => return Ok(()),
};
crate::operational_metrics::record_fts_build_attempt();
self.fts_build_attempts.fetch_add(1, Ordering::Relaxed);
table
.create_index(&["content"], Index::FTS(FtsIndexBuilder::default()))
.replace(true)
.execute()
.await
.context("Failed to (re)build FTS index")?;
crate::operational_metrics::record_fts_rebuild();
self.fts_indexed.store(true, Ordering::Relaxed);
Ok(())
}
}
impl VectorStore {
pub(super) async fn optimize_indices_locked(
&self,
tables: &[(&'static str, &Table)],
) -> Result<()> {
let mut first_err = None;
for (name, table) in tables {
let mut r = retry_on_conflict!(
table,
table.optimize(OptimizeAction::Index(
lancedb::table::OptimizeOptions::default()
))
)
.with_context(|| format!("Failed to optimize indices of {name} table"));
if should_rebuild_fts_after(name, &r) {
tracing::warn!(
error = %r.as_ref().expect_err("checked above"),
"FTS incremental optimize panicked; rebuilding the FTS index instead"
);
r = self
.rebuild_fts_index()
.await
.map(|()| lancedb::table::OptimizeStats::default())
.context("Failed to rebuild FTS index after incremental optimize panic");
}
if let Err(e) = r {
first_err.get_or_insert(e);
}
}
first_err.map_or(Ok(()), Err)
}
}
pub(super) fn should_rebuild_fts_after<T>(table: &str, r: &Result<T>) -> bool {
table == "chunks"
&& r.as_ref()
.err()
.is_some_and(|e| super::is_fts_compaction_panic(e))
}