use std::collections::{BTreeMap, BTreeSet};
use uqa_core::{DocId, Payload, PostingEntry, PostingList, Predicate, Value};
use uqa_storage::{BTreeIndex, DocumentStore};
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, String)>,
}
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
}
pub(crate) struct ColumnValueIndex {
index: BTreeIndex,
nulls: Vec<DocId>,
has_temporal: bool,
}
fn value_is_temporal(value: &Value) -> bool {
matches!(value, Value::Temporal(_))
}
fn value_is_nan(value: &Value) -> bool {
matches!(value, Value::Float(f) if f.is_nan())
}
fn predicate_targets_are_index_safe(predicate: &Predicate) -> bool {
let safe = |v: &Value| !value_is_temporal(v) && !value_is_nan(v);
match predicate {
Predicate::Equals(v)
| Predicate::NotEquals(v)
| Predicate::GreaterThan(v)
| Predicate::GreaterThanOrEqual(v)
| Predicate::LessThan(v)
| Predicate::LessThanOrEqual(v) => safe(v),
Predicate::InSet(values) => values.iter().all(safe),
Predicate::Between { low, high } => safe(low) && safe(high),
Predicate::IsNull | Predicate::IsNotNull => true,
}
}
impl ColumnValueIndex {
pub(crate) fn build(field: &str, values: impl Iterator<Item = (DocId, Value)>) -> Self {
let mut index = BTreeIndex::new(field);
let mut nulls = Vec::new();
let mut has_temporal = false;
for (doc_id, value) in values {
match value {
Value::Null => nulls.push(doc_id),
value => {
has_temporal |= value_is_temporal(&value);
index.insert(doc_id, value);
}
}
}
nulls.sort_unstable();
nulls.dedup();
Self {
index,
nulls,
has_temporal,
}
}
pub(crate) fn insert(&mut self, doc_id: DocId, value: &Value) {
match value {
Value::Null => {
if let Err(pos) = self.nulls.binary_search(&doc_id) {
self.nulls.insert(pos, doc_id);
}
}
value => {
self.has_temporal |= value_is_temporal(value);
self.index.insert(doc_id, value.clone());
}
}
}
pub(crate) fn remove(&mut self, doc_id: DocId, value: &Value) {
match value {
Value::Null => {
if let Ok(pos) = self.nulls.binary_search(&doc_id) {
self.nulls.remove(pos);
}
}
value => self.index.remove(doc_id, value),
}
}
pub(crate) fn clear(&mut self) {
self.index.clear();
self.nulls.clear();
self.has_temporal = false;
}
pub(crate) fn scan(&self, predicate: &Predicate) -> Option<PostingList> {
if !self.supports(predicate) {
return None;
}
match predicate {
Predicate::IsNull => Some(posting_list_from_sorted_ids(self.nulls.iter().copied())),
Predicate::IsNotNull => Some(self.index.scan(&Predicate::IsNotNull)),
Predicate::NotEquals(_) => unreachable!("unsupported predicates return above"),
predicate => Some(self.index.scan(predicate)),
}
}
pub(crate) fn estimate_cardinality(&self, predicate: &Predicate) -> Option<usize> {
if !self.supports(predicate) {
return None;
}
Some(match predicate {
Predicate::IsNull => self.nulls.len(),
Predicate::IsNotNull => self.index.estimate_cardinality(predicate),
Predicate::NotEquals(_) => unreachable!("unsupported predicates return above"),
predicate => self.index.estimate_cardinality(predicate),
})
}
fn supports(&self, predicate: &Predicate) -> bool {
predicate_targets_are_index_safe(predicate)
&& !matches!(predicate, Predicate::NotEquals(_))
&& (matches!(predicate, Predicate::IsNull | Predicate::IsNotNull) || !self.has_temporal)
}
}
fn posting_list_from_sorted_ids(ids: impl Iterator<Item = DocId>) -> PostingList {
let entries: Vec<PostingEntry> = ids
.map(|doc_id| PostingEntry::new(doc_id, Payload::default()))
.collect();
PostingList::from_sorted_unchecked(entries)
}
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_indexable_fields(&self, table: &str) -> StorageBackendResult<Vec<String>> {
let Some(table_name) = self.try_resolve_table_name(table)? else {
return Ok(Vec::new());
};
let mut fields = Vec::new();
if let Some(t) = self.try_table(&table_name)? {
for column in t.columns.read().iter() {
if (column.primary_key || column.unique) && !fields.contains(&column.name) {
fields.push(column.name.clone());
}
}
for constraint in t.key_constraints.read().iter() {
for column in &constraint.columns {
if !fields.contains(column) {
fields.push(column.clone());
}
}
}
}
for row in self.durable.catalog_indexes.read().values() {
if !row.index_type.eq_ignore_ascii_case("btree") {
continue;
}
if row.table_name != table_name {
continue;
}
let columns: Vec<String> = serde_json::from_str(&row.columns_json)?;
if let Some(first) = columns.first() {
if !fields.contains(first) {
fields.push(first.clone());
}
}
}
Ok(fields)
}
pub(crate) fn value_index_scan(
&self,
table: &str,
field: &str,
predicate: &Predicate,
) -> Result<Option<PostingList>, SQLError> {
let t = self.require_query_table(table)?;
{
let indexes = t.value_indexes.read();
if let Some(index) = indexes.get(field) {
return Ok(index.scan(predicate));
}
}
if !self
.ensure_query_value_index(table, &t, field)
.map_err(|error| SQLError::Internal(format!("build value index: {error}")))?
{
return Ok(None);
}
let result = t
.value_indexes
.read()
.get(field)
.and_then(|index| index.scan(predicate));
Ok(result)
}
pub(crate) fn value_index_cardinality(
&self,
table: &str,
field: &str,
predicate: &Predicate,
) -> Result<Option<usize>, SQLError> {
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 Some(table_name) = self.try_resolve_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: &str,
) -> 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);
}
}
if !self
.value_indexable_fields(table_name)?
.iter()
.any(|name| name == field)
{
return Ok(false);
}
let store = table.document_store.read();
let values =
Self::project_value_index_rows(store.as_ref(), table_name, field, store.doc_ids()?)?;
table.value_indexes.write().insert(
field.to_string(),
ColumnValueIndex::build(field, values.into_iter()),
);
Ok(true)
}
fn ensure_value_index(&self, table: &str, field: &str) -> StorageBackendResult<bool> {
self.ensure_value_index_with_mode(table, field, MissingValueIndexMode::MemoryOnly)
}
fn ensure_persistent_value_index(
&self,
table: &str,
field: &str,
) -> StorageBackendResult<bool> {
self.ensure_value_index_with_mode(table, field, MissingValueIndexMode::Persist)
}
fn project_value_index_rows(
store: &dyn DocumentStore,
table_name: &str,
field: &str,
doc_ids: Vec<DocId>,
) -> StorageBackendResult<Vec<(DocId, Value)>> {
let mut projected = store.get_fields_multi(&doc_ids, &[field])?;
let mut values = Vec::with_capacity(doc_ids.len());
for doc_id in doc_ids {
let Some(row) = projected.remove(&doc_id) else {
if store.get(doc_id)?.is_none() {
continue;
}
return Err(StorageBackendError::Other(format!(
"value-index rebuild for `{table_name}`.`{field}` lost document {doc_id} from the field projection"
)));
};
let [value]: [Value; 1] = row.try_into().map_err(|row: Vec<Value>| {
StorageBackendError::Other(format!(
"value-index rebuild for `{table_name}`.`{field}` returned {} projected values for document {doc_id}; expected 1",
row.len()
))
})?;
values.push((doc_id, value));
}
Ok(values)
}
fn ensure_value_index_with_mode(
&self,
table: &str,
field: &str,
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 store = t.document_store.read();
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 = store.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(store.as_ref(), &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(
store.as_ref(),
&table_name,
field,
store.doc_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, values.into_iter());
let mut indexes = t.value_indexes.write();
if support_changed {
indexes.insert(field.to_string(), built);
} else {
indexes.entry(field.to_string()).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<String> = 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(())
}
fn rebuild_persistent_value_indexes(
table_name: &str,
table: &TableState,
fields: &[String],
backend: &dyn uqa_storage::PersistentStorageBackend,
) -> StorageBackendResult<()> {
if fields.is_empty() {
return Ok(());
}
let field_refs = fields.iter().map(String::as_str).collect::<Vec<_>>();
let store = table.document_store.read();
let doc_ids = store.doc_ids()?;
let mut projected = store.get_fields_multi(&doc_ids, &field_refs)?;
let mut values_by_field = fields
.iter()
.map(|_| Vec::with_capacity(doc_ids.len()))
.collect::<Vec<Vec<(DocId, Value)>>>();
for doc_id in doc_ids {
let Some(values) = projected.remove(&doc_id) else {
if store.get(doc_id)?.is_none() {
continue;
}
return Err(StorageBackendError::Other(format!(
"value-index rebuild for `{table_name}` lost document {doc_id} from the field projection"
)));
};
if values.len() != fields.len() {
return Err(StorageBackendError::Other(format!(
"value-index rebuild for `{table_name}` returned {} projected values for document {doc_id}; expected {}",
values.len(),
fields.len()
)));
}
for (index, value) in values.into_iter().enumerate() {
values_by_field[index].push((doc_id, value));
}
}
let replacements = fields
.iter()
.zip(&values_by_field)
.map(|(field, values)| (field.as_str(), values.as_slice()))
.collect::<Vec<_>>();
backend.replace_btree_indexes(table_name, &replacements)?;
let built = fields
.iter()
.cloned()
.zip(values_by_field)
.map(|(field, values)| {
let index = ColumnValueIndex::build(&field, values.into_iter());
(field, index)
})
.collect::<Vec<_>>();
let mut indexes = table.value_indexes.write();
for (field, index) in built {
indexes.entry(field).or_insert(index);
}
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()? {
let desired: BTreeSet<String> =
self.value_indexable_fields(&table)?.into_iter().collect();
let actual: BTreeSet<String> =
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<()> {
if !self
.storage
.backend
.as_ref()
.is_some_and(|backend| backend.persists_btree_indexes())
{
return Ok(());
}
for table in self.table_names()? {
let Some(t) = self.try_table(&table)? else {
continue;
};
let fields: Vec<String> = t.value_indexes.read().keys().cloned().collect();
t.value_indexes.write().clear();
for field in fields {
self.ensure_value_index(&table, &field)?;
}
}
Ok(())
}
pub(crate) fn persistent_value_index_document_values(
&self,
table: &str,
document: &BTreeMap<String, Value>,
) -> Result<Option<BTreeMap<String, Value>>, SQLError> {
if !self
.storage
.backend
.as_ref()
.is_some_and(|backend| backend.persists_btree_indexes())
{
return Ok(None);
}
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(None);
}
let fields = self
.value_indexable_fields(&table_name)
.map_err(|err| SQLError::Internal(format!("read value-index policy: {err}")))?;
Ok(Some(
fields
.into_iter()
.map(|field| {
let value = document.get(&field).cloned().unwrap_or(Value::Null);
(field, value)
})
.collect(),
))
}
pub(crate) fn persist_value_indexes_apply_write(
&self,
table: &str,
doc_id: DocId,
new: Option<&BTreeMap<String, 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<String, Value>>,
new: Option<&BTreeMap<String, 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<String>> {
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,
) -> Result<Option<BTreeMap<String, Value>>, SQLError> {
let fields: Vec<String> = {
let indexes = t.value_indexes.read();
if indexes.is_empty() {
return Ok(None);
}
indexes.keys().cloned().collect()
};
let field_refs: Vec<&str> = fields.iter().map(String::as_str).collect();
let mut rows = t
.document_store
.read()
.get_fields_multi(&[doc_id], &field_refs)
.map_err(|error| SQLError::Internal(format!("read indexed fields: {error}")))?;
let values = rows
.remove(&doc_id)
.unwrap_or_else(|| vec![Value::Null; fields.len()]);
Ok(Some(fields.into_iter().zip(values).collect()))
}
pub(crate) fn value_indexes_clear(t: &TableState) {
t.value_indexes.write().clear();
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use uqa_storage::document_store::{Document, DocumentStore};
#[derive(Clone)]
struct MissingProjectionStore;
impl DocumentStore for MissingProjectionStore {
fn put(&mut self, _doc_id: DocId, _document: Document) -> StorageBackendResult<()> {
Ok(())
}
fn get(&self, doc_id: DocId) -> StorageBackendResult<Option<Document>> {
Ok((doc_id == 1).then(Document::new))
}
fn delete(&mut self, _doc_id: DocId) -> StorageBackendResult<()> {
Ok(())
}
fn clear(&mut self) -> StorageBackendResult<()> {
Ok(())
}
fn get_fields_multi(
&self,
_doc_ids: &[DocId],
_fields: &[&str],
) -> StorageBackendResult<BTreeMap<DocId, Vec<Value>>> {
Ok(BTreeMap::new())
}
fn doc_ids(&self) -> StorageBackendResult<Vec<DocId>> {
Ok(vec![1])
}
fn len(&self) -> StorageBackendResult<usize> {
Ok(1)
}
fn snapshot(&self) -> StorageBackendResult<Arc<dyn DocumentStore>> {
Ok(Arc::new(self.clone()))
}
fn writable_snapshot(&self) -> StorageBackendResult<Box<dyn DocumentStore>> {
Ok(Box::new(self.clone()))
}
}
fn ids(list: &PostingList) -> Vec<DocId> {
list.entries().iter().map(|e| e.doc_id).collect()
}
#[test]
fn build_scan_equals_and_ranges() {
let index = ColumnValueIndex::build(
"qty",
vec![
(1, Value::Int(10)),
(2, Value::Int(20)),
(3, Value::Int(20)),
(4, Value::Null),
(5, Value::Int(30)),
]
.into_iter(),
);
assert_eq!(
ids(&index.scan(&Predicate::Equals(Value::Int(20))).unwrap()),
vec![2, 3]
);
assert_eq!(
ids(&index.scan(&Predicate::GreaterThan(Value::Int(10))).unwrap()),
vec![2, 3, 5]
);
assert_eq!(
ids(&index
.scan(&Predicate::Between {
low: Value::Int(10),
high: Value::Int(20),
})
.unwrap()),
vec![1, 2, 3]
);
assert_eq!(ids(&index.scan(&Predicate::IsNull).unwrap()), vec![4]);
assert_eq!(
ids(&index.scan(&Predicate::IsNotNull).unwrap()),
vec![1, 2, 3, 5]
);
assert!(index.scan(&Predicate::NotEquals(Value::Int(10))).is_none());
}
#[test]
fn incremental_insert_remove_tracks_nulls() {
let mut index = ColumnValueIndex::build("qty", std::iter::empty());
index.insert(7, &Value::Int(1));
index.insert(8, &Value::Null);
assert_eq!(
ids(&index.scan(&Predicate::Equals(Value::Int(1))).unwrap()),
vec![7]
);
assert_eq!(ids(&index.scan(&Predicate::IsNull).unwrap()), vec![8]);
index.remove(7, &Value::Int(1));
index.remove(8, &Value::Null);
assert!(ids(&index.scan(&Predicate::Equals(Value::Int(1))).unwrap()).is_empty());
assert!(ids(&index.scan(&Predicate::IsNull).unwrap()).is_empty());
}
#[test]
fn temporal_and_nan_guards_refuse_acceleration() {
let temporal = uqa_core::TemporalValue::parse_date("2024-01-01").unwrap();
let index = ColumnValueIndex::build(
"ts",
vec![(1, Value::Temporal(temporal.clone()))].into_iter(),
);
assert!(index
.scan(&Predicate::Equals(Value::Str("2024-01-01".into())))
.is_none());
let numeric = ColumnValueIndex::build("f", vec![(1, Value::Float(1.0))].into_iter());
assert!(numeric
.scan(&Predicate::Equals(Value::Float(f64::NAN)))
.is_none());
assert!(numeric
.scan(&Predicate::Equals(Value::Temporal(temporal)))
.is_none());
}
#[test]
fn rebuild_rejects_a_document_missing_from_the_field_projection() {
let engine = crate::Engine::new();
engine
.sql("CREATE TABLE projection_gap (id INTEGER PRIMARY KEY)", &[])
.unwrap();
let table = engine.try_table("projection_gap").unwrap().unwrap();
*table.document_store.write() = Box::new(MissingProjectionStore);
crate::Engine::value_indexes_clear(&table);
let error = engine
.ensure_value_index("projection_gap", "id")
.unwrap_err();
assert!(error.to_string().contains("lost document 1"), "{error}");
assert!(table.value_indexes.read().is_empty());
}
#[test]
fn relation_key_suffix_preserves_quoted_components() {
assert_eq!(unqualified_relation_key("public.items"), Some("items"));
assert_eq!(
unqualified_relation_key("public.\"items.with.dot\""),
Some("\"items.with.dot\"")
);
assert_eq!(
unqualified_relation_key("\"schema.with.dot\".\"items.with.dot\""),
Some("\"items.with.dot\"")
);
assert_eq!(
unqualified_relation_key("public.\"items\"\"quoted\""),
Some("\"items\"\"quoted\"")
);
}
#[test]
fn query_builds_missing_durable_index_in_memory_only() {
let directory = tempfile::tempdir().unwrap();
let engine = crate::Engine::open(&directory.path().join("memory-only-btree.db")).unwrap();
engine
.sql("CREATE TABLE items (id INTEGER PRIMARY KEY)", &[])
.unwrap();
engine
.sql("INSERT INTO items (id) VALUES (1)", &[])
.unwrap();
let backend = engine.storage.backend.as_ref().unwrap();
backend.drop_btree_index("public.items", "id").unwrap();
let table = engine.try_table("items").unwrap().unwrap();
crate::Engine::value_indexes_clear(&table);
let result = engine
.sql("SELECT id FROM items WHERE id = 1", &[])
.unwrap();
assert_eq!(result.rows.len(), 1);
assert!(backend
.load_btree_index("public.items", "id")
.unwrap()
.is_none());
assert!(engine
.try_table("items")
.unwrap()
.unwrap()
.value_indexes
.read()
.contains_key("id"));
engine.reload_persistent_value_indexes().unwrap();
assert!(backend
.load_btree_index("public.items", "id")
.unwrap()
.is_none());
assert!(engine
.try_table("items")
.unwrap()
.unwrap()
.value_indexes
.read()
.contains_key("id"));
engine.ensure_persistent_value_index("items", "id").unwrap();
assert_eq!(
backend
.load_btree_index("public.items", "id")
.unwrap()
.unwrap(),
vec![(1, Value::Int(1))]
);
}
#[test]
fn open_repair_discards_raw_alias_and_rebuilds_canonical_index() {
let directory = tempfile::tempdir().unwrap();
let database = directory.path().join("repair-btree.db");
let engine = crate::Engine::open(&database).unwrap();
engine
.sql("CREATE TABLE items (id INTEGER PRIMARY KEY)", &[])
.unwrap();
engine
.sql("INSERT INTO items (id) VALUES (1)", &[])
.unwrap();
let backend = engine.storage.backend.as_ref().unwrap().clone();
backend.drop_btree_index("public.items", "id").unwrap();
backend
.replace_btree_index("public.items", "obsolete", &[(1, Value::Int(888))])
.unwrap();
let raw = rusqlite::Connection::open(&database).unwrap();
raw.execute("DROP TRIGGER _btree_entries_document_insert", [])
.unwrap();
backend
.replace_btree_index("items", "id", &[(1, Value::Int(999))])
.unwrap();
raw.execute_batch(
"CREATE TRIGGER _btree_entries_document_insert
BEFORE INSERT ON _btree_index_entries
WHEN NOT EXISTS (
SELECT 1 FROM _documents
WHERE table_name = NEW.table_name AND doc_id = NEW.doc_id
)
BEGIN
SELECT RAISE(ABORT, 'persistent B-tree entry has no backing document');
END;",
)
.unwrap();
drop(raw);
let table = engine.try_table("items").unwrap().unwrap();
crate::Engine::value_indexes_clear(&table);
drop(table);
drop(backend);
drop(engine);
let reopened = crate::Engine::open(&database).unwrap();
let backend = reopened.storage.backend.as_ref().unwrap();
assert!(backend.load_btree_index("items", "id").unwrap().is_none());
assert!(backend
.load_btree_index("public.items", "obsolete")
.unwrap()
.is_none());
assert_eq!(
backend
.load_btree_index("public.items", "id")
.unwrap()
.unwrap(),
vec![(1, Value::Int(1))]
);
}
#[test]
fn clean_open_repair_does_not_contend_for_sqlite_writer_lock() {
let directory = tempfile::tempdir().unwrap();
let database = directory.path().join("clean-repair.db");
let engine = crate::Engine::open(&database).unwrap();
engine
.sql("CREATE TABLE items (id INTEGER PRIMARY KEY)", &[])
.unwrap();
engine
.sql("INSERT INTO items (id) VALUES (1)", &[])
.unwrap();
assert!(engine
.persistent_value_index_repair_plan()
.unwrap()
.is_empty());
let blocker = engine
.storage
.provider
.as_ref()
.unwrap()
.open_session()
.unwrap();
blocker.backend.begin_transaction().unwrap();
let repair_result = engine.repair_persistent_value_indexes_on_open();
let new_session_result = engine.new_session();
let reopen_result = crate::Engine::open(&database);
blocker.backend.rollback_transaction().unwrap();
repair_result.unwrap();
new_session_result.unwrap();
reopen_result.unwrap();
}
}