//
// Unified Query Algebra
//
// Copyright (c) 2023-2026 Cognica, Inc.
//
//! Retain value-index caches and provider handles for execution-owned index selection and lookup.
//!
//! Catalog policy selects column accelerators and named expression indexes. Query hydration is memory-only; DDL and repair may publish durable postings. Execution owns predicate eligibility, NULL handling, stored-key maintenance and result construction through [`ColumnValueIndex`].
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 {
/// Build an in-memory accelerator from the pinned document snapshot, but
/// leave durable storage untouched. Query execution and rollback recovery
/// use this mode so a read transaction can never be upgraded by an index
/// cache miss.
MemoryOnly,
/// Materialize the complete durable posting set when it is absent. Only
/// DDL and the explicit open-time repair boundary may use this mode.
Persist,
}
#[derive(Debug, Default, PartialEq, Eq)]
struct PersistentValueIndexRepairPlan {
/// Legacy unqualified table keys whose complete durable posting sets must
/// be removed before their canonical counterparts are rebuilt.
aliases: BTreeSet<String>,
/// Canonical tables whose durable marker fields differ from catalog policy
/// or which had a legacy alias.
tables: BTreeSet<String>,
/// Durable retry markers written by a catalog migration. They are cleared
/// in the same transaction, and only after every requested repair succeeds.
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),
}
}
/// Estimate one exact value-index predicate without materializing or
/// sorting its posting list. Engine column indexes keep every document in
/// one value bucket, so the storage upper bound is exact here.
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)
}
/// Return whether catalog policy provides an exact in-memory value-index
/// implementation for this predicate. Missing hot state is hydrated in
/// memory, preserving the read-only lazy-recovery contract without forcing
/// the relational planner to execute every scalar filter as a posting scan.
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)
}
/// Make the accelerators of `fields` available to an index-only read of `table`, the handle a query bound for `name`. Only the live table loads or builds an accelerator: a detached snapshot table lives for one statement, and building its accelerators would read every document to save reading a few.
pub(crate) fn prepare_index_only_read(
&self,
name: &str,
table: &std::sync::Arc<dyn uqa_execution::query::table_read::TableRead>,
fields: &[String],
) -> Result<bool, SQLError> {
if !self.index_only_scans_enabled() {
return Ok(false);
}
let prepare = || -> StorageBackendResult<bool> {
let Some(table_name) = self.try_resolve_query_table_name(name)? else {
return Ok(false);
};
let Some(live) = self.try_table(&table_name)? else {
return Ok(false);
};
if !std::ptr::addr_eq(std::sync::Arc::as_ptr(&live), std::sync::Arc::as_ptr(table)) {
return Ok(false);
}
for field in fields {
let field = ValueIndexKey::Column(field.clone());
// A loaded accelerator answers without consulting the catalog or the durable postings again.
if !self.ensure_query_value_index(&table_name, &live, &field)? {
return Ok(false);
}
}
// A read of no field asks whichever accelerator its access path has loaded by then whether each row exists.
Ok(true)
};
prepare().map_err(|error| {
uqa_execution::storage_errors::storage_error("prepare index-only read", &error)
})
}
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)?;
let built = self.build_value_index(&table_name, table, field, values)?;
table.value_indexes.write().insert(field.clone(), built);
Ok(true)
}
/// Hydrate one value index from durable postings when available. A missing
/// durable marker is satisfied by an in-memory build only; query execution
/// must not turn a deferred read transaction into a writer.
fn ensure_value_index(&self, table: &str, field: &ValueIndexKey) -> StorageBackendResult<bool> {
self.ensure_value_index_with_mode(table, field, MissingValueIndexMode::MemoryOnly)
}
/// DDL/open-repair counterpart of [`Engine::ensure_value_index`].
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 {
// Keep every posting that still has an authoritative document
// and parse only documents whose posting is missing. Historical
// inconsistencies are normally sparse; rebuilding the complete
// field could otherwise parse gigabytes to repair one row.
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 = self.build_value_index(&table_name, &t, field, values)?;
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)
}
/// Reconcile one table's in-memory and durable indexes with its current
/// PRIMARY KEY / UNIQUE / catalog-btree policy.
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);
}
}
let carried = self.value_index_carried_fields(&table_name, &t)?;
{
let mut indexes = t.value_indexes.write();
indexes.retain(|field, _| desired.contains(field));
// An index definition may have turned a carried column into a search key, or the reverse.
for (field, index) in indexes.iter_mut() {
let carry = carried.contains(field);
if index.is_carried() != carry {
let current = std::mem::replace(
index,
ColumnValueIndex::build_carried(std::iter::empty()),
);
*index = current.with_use(field.name(), carry);
}
}
}
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(())
}
/// Reconcile durable value indexes at the explicit database-open repair
/// boundary. A read-only preflight keeps the normal open/session path out
/// of `SQLite`'s single-writer lane. Only an observed missing/stale marker,
/// pending structural repair, or pre-canonicalization alias opens the writer
/// transaction, where the plan is recomputed against the pinned snapshot
/// before making any changes.
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| {
// Waiting for the writer reservation may have made the preflight
// stale. Recompute after the transaction has refreshed its pinned
// catalog/data snapshot and mutate only what is still divergent.
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)
}
/// Restore hot accelerators directly from rolled-back postings. Recovery can hold the transaction mutex, so it must never bind SQL expressions or execute callbacks; missing indexes remain cold until the next statement.
/// Drop the in-memory value indexes of every durable table after a rollback. The table catalog reload that follows every rollback builds a new state for each durable table, without indexes, so reloading them from storage beforehand was thrown away: a query builds the indexes it needs again from the snapshot it reads. Dropping them here keeps the earlier states from serving rolled-back keys when that reload fails. A temporary table's indexes are restored with its data snapshot instead, and a memory-only engine restores every table that way.
pub(crate) fn drop_persistent_value_indexes(&self) {
if self.storage.backend.is_none() {
return;
}
let tables = self
.storage
.tables
.read()
.values()
.cloned()
.collect::<Vec<_>>();
for table in tables {
if table.persistence != uqa_sql::ast::RelationPersistence::Temporary {
table.value_indexes.write().clear();
}
}
}
/// `unused` names the namespace in which no document ever had `doc_id`, so that the document's entries replace none.
pub(crate) fn persist_value_indexes_apply_write(
&self,
table: &str,
doc_id: DocId,
new: Option<&BTreeMap<ValueIndexKey, Value>>,
unused: Option<uqa_storage::document_store::identifiers::DocumentIdNamespace>,
) -> 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(());
}
match (new, unused) {
(Some(new), Some(namespace)) => {
backend.apply_unused_btree_index_write(&table_name, doc_id, new, namespace)
}
(new, _) => backend.apply_btree_index_write(&table_name, doc_id, new),
}
.map_err(|err| SQLError::Internal(format!("btree index write failed: {err}")))
}
/// TRUNCATE keeps index definitions installed but removes all postings.
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(())
}
/// Incremental maintenance for built indexes. `old` carries the
/// previous field values when the document already existed.
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));
}
}
}
/// Names of every built value-index field, or `None` when no index
/// is built. Known-new writes use this instead of
/// [`Engine::value_indexes_old_values`], because a document id that
/// was never stored has no previous values worth a storage lookup.
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())
}
/// Read the actual cached keys, which may differ from re-evaluating a replaced immutable function on the old document.
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()
})
}
/// Drop every built index for the table (TRUNCATE, bulk reloads,
/// store replacement, schema changes).
pub(crate) fn value_indexes_clear(t: &TableState) {
t.value_indexes.write().clear();
}
}
#[cfg(test)]
mod tests;