use std::collections::{BTreeMap, BTreeSet};
use uqa_core::{DocId, PostingList, Predicate, Value};
pub(crate) use uqa_execution::catalog::index::value::ColumnValueIndex;
use uqa_storage::ValueIndexKey;
mod keys;
use crate::{SQLError, StorageBackendError, StorageBackendResult, TableState};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum MissingValueIndexMode {
MemoryOnly,
Persist,
}
#[derive(Debug, Default, PartialEq, Eq)]
struct PersistentValueIndexRepairPlan {
aliases: BTreeSet<String>,
tables: BTreeSet<String>,
pending: BTreeSet<(String, ValueIndexKey)>,
}
impl PersistentValueIndexRepairPlan {
fn is_empty(&self) -> bool {
self.aliases.is_empty() && self.tables.is_empty() && self.pending.is_empty()
}
}
fn unqualified_relation_key(qualified: &str) -> Option<&str> {
let mut quoted = false;
let mut chars = qualified.char_indices().peekable();
while let Some((index, ch)) = chars.next() {
if ch == '"' {
if quoted && chars.peek().is_some_and(|(_, next)| *next == '"') {
chars.next();
} else {
quoted = !quoted;
}
} else if ch == '.' && !quoted {
return Some(&qualified[index + 1..]);
}
}
None
}
impl crate::Engine {
fn persistent_value_index_backend(
&self,
table: &TableState,
) -> Option<&dyn uqa_storage::PersistentStorageBackend> {
if table.persistence == uqa_sql::ast::RelationPersistence::Temporary {
return None;
}
self.storage
.backend
.as_deref()
.filter(|backend| backend.persists_btree_indexes())
}
fn value_index_table_is_temporary(&self, table: &str) -> Result<bool, SQLError> {
self.try_table(table)
.map_err(|err| SQLError::Internal(format!("resolve value-index table: {err}")))?
.map(|table| table.persistence == uqa_sql::ast::RelationPersistence::Temporary)
.ok_or_else(|| SQLError::UnknownTable(table.to_string()))
}
pub(crate) fn value_index_scan_key(
&self,
table: &str,
field: &ValueIndexKey,
predicate: &Predicate,
) -> Result<Option<PostingList>, SQLError> {
let state = self.require_table(table)?;
self.value_index_scan_state(table, &state, field, predicate, None)
}
pub(crate) fn value_index_query_scan(
&self,
table: &str,
field: &str,
predicate: &Predicate,
) -> Result<Option<PostingList>, SQLError> {
let read = self.serializable_table_read(table)?;
let state = self.require_query_table(table)?;
self.value_index_scan_state(
table,
&state,
&ValueIndexKey::Column(field.into()),
predicate,
read.as_ref(),
)
}
pub(crate) fn value_index_scan_state(
&self,
table: &str,
t: &std::sync::Arc<TableState>,
field: &ValueIndexKey,
predicate: &Predicate,
read: Option<&uqa_execution::serializable::SerializableRelationRead>,
) -> Result<Option<PostingList>, SQLError> {
let observed = read.map(|read| (read, t.columns.snapshot()));
let scan = |index: &ColumnValueIndex| {
index.scan_observing(predicate, || {
if let Some((read, columns)) = &observed {
read.observe_column_index(columns, field, predicate)?;
}
Ok(())
})
};
{
let indexes = t.value_indexes.read();
if let Some(index) = indexes.get(field) {
return scan(index);
}
}
if !self
.ensure_query_value_index(table, t, field)
.map_err(|error| {
uqa_execution::storage_errors::storage_error("build value index", &error)
})?
{
return Ok(None);
}
let indexes = t.value_indexes.read();
match indexes.get(field) {
Some(index) => scan(index),
None => Ok(None),
}
}
pub(crate) fn value_index_cardinality(
&self,
table: &str,
field: &str,
predicate: &Predicate,
) -> Result<Option<usize>, SQLError> {
let field = &ValueIndexKey::Column(field.into());
let table_state = self.require_query_table(table)?;
{
let indexes = table_state.value_indexes.read();
if let Some(index) = indexes.get(field) {
return Ok(index.estimate_cardinality(predicate));
}
}
if !self
.ensure_query_value_index(table, &table_state, field)
.map_err(|error| SQLError::Internal(format!("build value index: {error}")))?
{
return Ok(None);
}
let cardinality = table_state
.value_indexes
.read()
.get(field)
.and_then(|index| index.estimate_cardinality(predicate));
Ok(cardinality)
}
pub(crate) fn value_index_supports(
&self,
table: &str,
field: &str,
predicate: &Predicate,
) -> StorageBackendResult<bool> {
let field = &ValueIndexKey::Column(field.into());
let Some(table_name) = self.try_resolve_query_table_name(table)? else {
return Ok(false);
};
let Some(table) = self.try_query_table(&table_name)? else {
return Ok(false);
};
if !self.ensure_query_value_index(&table_name, &table, field)? {
return Ok(false);
}
let supported = table
.value_indexes
.read()
.get(field)
.is_some_and(|index| index.supports(predicate));
Ok(supported)
}
fn ensure_query_value_index(
&self,
table_name: &str,
table: &std::sync::Arc<TableState>,
field: &ValueIndexKey,
) -> StorageBackendResult<bool> {
if table.value_indexes.read().contains_key(field) {
return Ok(true);
}
if let Some(live) = self.try_table(table_name)? {
if std::sync::Arc::ptr_eq(&live, table) {
return self.ensure_value_index(table_name, field);
}
}
let Some(table_name) = self.try_resolve_query_table_name(table_name)? else {
return Ok(false);
};
if !self
.value_indexable_fields_in_state(&table_name, table)?
.iter()
.any(|name| name == field)
{
return Ok(false);
}
let ids = table.document_store.read().doc_ids()?;
let values = self.project_value_index_rows(table, &table_name, field, &ids)?;
table.value_indexes.write().insert(
field.clone(),
ColumnValueIndex::build(field.name(), values.into_iter()),
);
Ok(true)
}
fn ensure_value_index(&self, table: &str, field: &ValueIndexKey) -> StorageBackendResult<bool> {
self.ensure_value_index_with_mode(table, field, MissingValueIndexMode::MemoryOnly)
}
fn ensure_persistent_value_index(
&self,
table: &str,
field: &ValueIndexKey,
) -> StorageBackendResult<bool> {
self.ensure_value_index_with_mode(table, field, MissingValueIndexMode::Persist)
}
fn ensure_value_index_with_mode(
&self,
table: &str,
field: &ValueIndexKey,
mode: MissingValueIndexMode,
) -> StorageBackendResult<bool> {
let Some(table_name) = self.try_resolve_table_name(table)? else {
return Ok(false);
};
let Some(t) = self.try_table(&table_name)? else {
return Ok(false);
};
let memory_index_exists = t.value_indexes.read().contains_key(field);
if !self
.value_indexable_fields(&table_name)?
.iter()
.any(|name| name == field)
{
return Ok(false);
}
let persistent_backend = self.persistent_value_index_backend(&t);
let persisted = persistent_backend
.map(|backend| backend.load_btree_index(&table_name, field))
.transpose()?
.flatten();
let durable_index_missing = persistent_backend.is_some() && persisted.is_none();
if memory_index_exists
&& (mode == MissingValueIndexMode::MemoryOnly || !durable_index_missing)
{
return Ok(true);
}
let (values, support_changed, repair_delta) = if let Some(values) = persisted {
let mut persisted_ids = values.iter().map(|(doc_id, _)| *doc_id).collect::<Vec<_>>();
persisted_ids.sort_unstable();
let mut document_ids = t.document_store.read().doc_ids()?;
document_ids.sort_unstable();
if persisted_ids == document_ids {
(values, false, None)
} else {
let document_id_set = document_ids.iter().copied().collect::<BTreeSet<_>>();
let mut present = BTreeSet::new();
let mut repaired = Vec::with_capacity(document_ids.len());
let mut stale = Vec::new();
for (doc_id, value) in values {
if document_id_set.contains(&doc_id) {
present.insert(doc_id);
repaired.push((doc_id, value));
} else {
stale.push(doc_id);
}
}
let missing = document_ids
.into_iter()
.filter(|doc_id| !present.contains(doc_id))
.collect::<Vec<_>>();
let missing = self.project_value_index_rows(&t, &table_name, field, &missing)?;
repaired.extend(missing.iter().cloned());
repaired.sort_unstable_by_key(|(doc_id, _)| *doc_id);
(repaired, true, Some((stale, missing)))
}
} else {
(
self.project_value_index_rows(&t, &table_name, field, &{
let ids = t.document_store.read().doc_ids()?;
ids
})?,
true,
None,
)
};
if support_changed && mode == MissingValueIndexMode::Persist {
if let Some(backend) = persistent_backend {
if let Some((stale, missing)) = repair_delta.as_ref() {
backend.repair_btree_index(&table_name, field, &values, stale, missing)?;
} else {
backend.replace_btree_index(&table_name, field, &values)?;
}
}
}
if !memory_index_exists || support_changed {
let built = ColumnValueIndex::build(field.name(), values.into_iter());
let mut indexes = t.value_indexes.write();
if support_changed {
indexes.insert(field.clone(), built);
} else {
indexes.entry(field.clone()).or_insert(built);
}
}
Ok(true)
}
pub(crate) fn refresh_value_indexes_for_table(&self, table: &str) -> StorageBackendResult<()> {
let table_name = self
.try_resolve_table_name(table)?
.ok_or_else(|| StorageBackendError::Other(format!("table `{table}` does not exist")))?;
let t = self.try_table(&table_name)?.ok_or_else(|| {
StorageBackendError::Other(format!("table `{table_name}` does not exist"))
})?;
let desired = self.value_indexable_fields(&table_name)?;
let mut stale: Vec<ValueIndexKey> = t
.value_indexes
.read()
.keys()
.filter(|field| !desired.contains(field))
.cloned()
.collect();
let persistent_backend = self.persistent_value_index_backend(&t);
let mut persisted_fields = BTreeSet::new();
if let Some(backend) = persistent_backend {
for field in backend.btree_index_fields(&table_name)? {
if !desired.contains(&field) && !stale.contains(&field) {
stale.push(field);
} else {
persisted_fields.insert(field);
}
}
for field in &stale {
backend.drop_btree_index(&table_name, field)?;
persisted_fields.remove(field);
}
}
t.value_indexes
.write()
.retain(|field, _| desired.contains(field));
if let Some(backend) = persistent_backend {
let missing = desired
.iter()
.filter(|field| !persisted_fields.contains(*field))
.cloned()
.collect::<Vec<_>>();
self.rebuild_persistent_value_indexes(&table_name, &t, &missing, backend)?;
for field in desired
.iter()
.filter(|field| persisted_fields.contains(*field))
{
self.ensure_persistent_value_index(&table_name, field)?;
}
} else {
for field in desired {
self.ensure_persistent_value_index(&table_name, &field)?;
}
}
Ok(())
}
pub(crate) fn repair_persistent_value_indexes_on_open(&self) -> StorageBackendResult<()> {
if self.persistent_value_index_repair_plan()?.is_empty() {
return Ok(());
}
self.with_implicit_storage_transaction(|engine| {
let plan = engine.persistent_value_index_repair_plan()?;
let Some(backend) = engine
.storage
.backend
.as_ref()
.filter(|backend| backend.persists_btree_indexes())
else {
return Ok(());
};
for alias in &plan.aliases {
for field in backend.btree_index_fields(alias)? {
backend.drop_btree_index(alias, &field)?;
}
}
for table in &plan.tables {
engine.refresh_value_indexes_for_table(table)?;
}
for (table, field) in &plan.pending {
if !plan.tables.contains(table) {
let should_exist = engine.try_table(table)?.is_some()
&& engine
.value_indexable_fields(table)?
.iter()
.any(|candidate| candidate == field);
if should_exist {
engine.ensure_persistent_value_index(table, field)?;
} else {
backend.drop_btree_index(table, field)?;
}
}
backend.clear_btree_index_repair(table, field)?;
}
Ok(())
})
}
fn persistent_value_index_repair_plan(
&self,
) -> StorageBackendResult<PersistentValueIndexRepairPlan> {
let Some(backend) = self
.storage
.backend
.as_ref()
.filter(|backend| backend.persists_btree_indexes())
else {
return Ok(PersistentValueIndexRepairPlan::default());
};
let mut plan = PersistentValueIndexRepairPlan {
pending: backend.btree_index_repairs()?.into_iter().collect(),
..PersistentValueIndexRepairPlan::default()
};
for table in self.table_names_in_execution()? {
let desired: BTreeSet<ValueIndexKey> =
self.value_indexable_fields(&table)?.into_iter().collect();
let actual: BTreeSet<ValueIndexKey> =
backend.btree_index_fields(&table)?.into_iter().collect();
let mut has_legacy_alias = false;
if let Some(alias) = unqualified_relation_key(&table) {
if !backend.btree_index_fields(alias)?.is_empty() {
has_legacy_alias = true;
plan.aliases.insert(alias.to_string());
}
}
if actual != desired || has_legacy_alias {
plan.tables.insert(table);
}
}
Ok(plan)
}
pub(crate) fn reload_persistent_value_indexes(&self) -> StorageBackendResult<()> {
let Some(backend) = self
.storage
.backend
.as_ref()
.filter(|backend| backend.persists_btree_indexes())
else {
return Ok(());
};
let tables = self
.storage
.tables
.read()
.iter()
.map(|(name, table)| (name.qualified_name(), table.clone()))
.collect::<Vec<_>>();
for (name, table) in tables {
if table.persistence == uqa_sql::ast::RelationPersistence::Temporary {
continue;
}
let fields = table
.value_indexes
.read()
.keys()
.cloned()
.collect::<Vec<_>>();
table.value_indexes.write().clear();
for field in fields {
if let Some(values) = backend.load_btree_index(&name, &field)? {
let index = ColumnValueIndex::build(field.name(), values.into_iter());
table.value_indexes.write().insert(field, index);
}
}
}
Ok(())
}
pub(crate) fn persist_value_indexes_apply_write(
&self,
table: &str,
doc_id: DocId,
new: Option<&BTreeMap<ValueIndexKey, Value>>,
) -> Result<(), SQLError> {
let Some(backend) = self
.storage
.backend
.as_ref()
.filter(|backend| backend.persists_btree_indexes())
else {
return Ok(());
};
let table_name = self
.try_resolve_table_name(table)
.map_err(|err| SQLError::Internal(format!("resolve value-index table: {err}")))?
.ok_or_else(|| SQLError::UnknownTable(table.to_string()))?;
if self.value_index_table_is_temporary(&table_name)? {
return Ok(());
}
backend
.apply_btree_index_write(&table_name, doc_id, new)
.map_err(|err| SQLError::Internal(format!("btree index write failed: {err}")))
}
pub(crate) fn value_indexes_truncate(
&self,
table: &str,
t: &TableState,
) -> Result<(), SQLError> {
let table_name = self
.try_resolve_table_name(table)
.map_err(|err| SQLError::Internal(format!("resolve value-index table: {err}")))?
.ok_or_else(|| SQLError::UnknownTable(table.to_string()))?;
if let Some(backend) = self.persistent_value_index_backend(t) {
backend
.clear_btree_indexes(&table_name)
.map_err(|err| SQLError::Internal(format!("btree truncate failed: {err}")))?;
}
for index in t.value_indexes.write().values_mut() {
index.clear();
}
Ok(())
}
pub(crate) fn value_indexes_apply_write(
t: &TableState,
doc_id: DocId,
old: Option<&BTreeMap<ValueIndexKey, Value>>,
new: Option<&BTreeMap<ValueIndexKey, Value>>,
) {
let mut indexes = t.value_indexes.write();
if indexes.is_empty() {
return;
}
for (field, index) in indexes.iter_mut() {
if let Some(old_values) = old {
index.remove(doc_id, old_values.get(field).unwrap_or(&Value::Null));
}
if let Some(new_values) = new {
index.insert(doc_id, new_values.get(field).unwrap_or(&Value::Null));
}
}
}
pub(crate) fn value_indexes_built_fields(t: &TableState) -> Option<Vec<ValueIndexKey>> {
let indexes = t.value_indexes.read();
if indexes.is_empty() {
return None;
}
Some(indexes.keys().cloned().collect())
}
pub(crate) fn value_indexes_old_values(
t: &TableState,
doc_id: DocId,
) -> Option<BTreeMap<ValueIndexKey, Value>> {
let indexes = t.value_indexes.read();
(!indexes.is_empty()).then(|| {
indexes
.iter()
.map(|(field, index)| {
(
field.clone(),
index.stored_value(doc_id).cloned().unwrap_or(Value::Null),
)
})
.collect()
})
}
pub(crate) fn value_indexes_clear(t: &TableState) {
t.value_indexes.write().clear();
}
pub(crate) fn value_indexes_clear_column_accelerators(t: &TableState) {
t.value_indexes
.write()
.retain(|key, _| matches!(key, ValueIndexKey::Index(_)));
}
}
#[cfg(test)]
mod tests;