mod sampling;
#[cfg(test)]
mod tests;
use std::sync::{atomic::Ordering, Arc};
use uqa_sql::ast::{ColumnDef, ColumnType, GeneratedColumnKind};
use uqa_storage::{StorageBackendError, StorageBackendResult};
use crate::statistics::{now_ms, MaintenanceState};
use crate::{ColumnStatsMap, Engine};
type DataGenerations = Vec<(String, Option<u64>)>;
struct AutomaticAnalysis {
object_id: [u8; 16],
maintenance: MaintenanceState,
columns: Arc<Vec<ColumnDef>>,
publication: Option<(u64, u64, Option<u64>)>,
data_generations: Option<DataGenerations>,
sampled_at: Option<u64>,
row_count: u64,
statistics: ColumnStatsMap,
}
fn automatic_column(ty: &ColumnType) -> bool {
match ty {
ColumnType::Domain { base, .. } => automatic_column(base),
ColumnType::Bytea
| ColumnType::Vector(_)
| ColumnType::Tensor(_)
| ColumnType::Json
| ColumnType::JsonB
| ColumnType::Array(_)
| ColumnType::AnyArray
| ColumnType::Record => false,
_ => true,
}
}
impl Engine {
pub(crate) fn run_automatic_analyze(&self, name: &str) -> StorageBackendResult<bool> {
let _statement = self.lock_statement_gate();
if self.storage.backend.is_none() {
return Ok(false);
}
self.with_storage_maintenance_scope(|engine| {
let Some(target) = uqa_execution::maintenance::analyze::prepare_optional_target(
&engine.analyze_execution_context(false),
name,
)
.map_err(|error| StorageBackendError::backend("automatic ANALYZE locking", error))?
else {
return Ok(false);
};
let Some(analysis) = engine.collect_automatic_analysis(&target.name)? else {
return Ok(false);
};
engine.publish_automatic_analysis(&target.name, analysis)
})
}
fn collect_automatic_analysis(
&self,
name: &str,
) -> StorageBackendResult<Option<AutomaticAnalysis>> {
let Some(table) = self.try_table(name)? else {
return Ok(None);
};
let Some(catalog) = self.storage.catalog.as_deref() else {
return Ok(None);
};
let maintenance = MaintenanceState::load_for(catalog, name, table.object_id())?;
let missing = maintenance.missing(table.column_stats.read().is_empty());
if !maintenance.due(
missing,
now_ms(),
crate::statistics::value_size::FORMAT_VERSION,
) {
return Ok(None);
}
let definitions = table.columns.snapshot();
let columns = definitions
.iter()
.filter(|column| {
automatic_column(&column.ty)
&& !column
.generated
.as_ref()
.is_some_and(|generated| generated.kind == GeneratedColumnKind::Virtual)
})
.map(|column| column.name.clone())
.collect::<Vec<_>>();
let data_generations = self.sampled_data_generations(name)?;
let publication = self.statistics_publication_identity(name)?;
let sampled_at = self.statistics_sample_sequence();
let (statistics, row_count) = sampling::collect(self, name, &columns)?;
Ok(Some(AutomaticAnalysis {
object_id: table.object_id(),
maintenance,
columns: definitions,
publication,
data_generations,
sampled_at,
row_count,
statistics,
}))
}
fn statistics_publication_identity(
&self,
name: &str,
) -> StorageBackendResult<Option<(u64, u64, Option<u64>)>> {
let Some(catalog) = self.storage.catalog.as_deref() else {
return Ok(None);
};
Ok(catalog.cache_revisions()?.map(|revisions| {
(
revisions.table_catalog,
revisions.storage_schema,
revisions.column_statistics.get(name).copied(),
)
}))
}
fn sampled_data_generations(
&self,
name: &str,
) -> StorageBackendResult<Option<DataGenerations>> {
let Some(catalog) = self.storage.catalog.as_deref() else {
return Ok(None);
};
let Some(revisions) = catalog.cache_revisions()? else {
return Ok(None);
};
let members = self
.hierarchy_scan_tables(name, true)
.map_err(|error| StorageBackendError::Other(error.to_string()))?;
Ok(Some(
members
.into_iter()
.map(|member| {
let generation = revisions.table_data.get(&member).copied();
(member, generation)
})
.collect(),
))
}
fn publish_automatic_analysis(
&self,
name: &str,
analysis: AutomaticAnalysis,
) -> StorageBackendResult<bool> {
self.with_storage_maintenance_scope(|engine| {
let Some(target) = uqa_execution::maintenance::analyze::prepare_optional_target(
&engine.analyze_execution_context(false),
name,
)
.map_err(|error| StorageBackendError::backend("automatic ANALYZE locking", error))?
else {
return Ok(false);
};
if target.object_id != analysis.object_id {
return Ok(false);
}
let name = target.name.as_str();
if !engine.try_prepare_storage_maintenance_writer()? {
return Ok(false);
}
let Some(table) = engine.try_table(name)? else {
return Ok(false);
};
let Some(catalog) = engine.storage.catalog.as_deref() else {
return Ok(false);
};
let current = MaintenanceState::load_for(catalog, name, table.object_id())?;
let columns = table.columns.snapshot();
let same_columns = columns.len() == analysis.columns.len()
&& columns.iter().zip(analysis.columns.iter()).all(|(a, b)| {
a.object_id == b.object_id
&& a.attribute_number == b.attribute_number
&& a.name == b.name
&& a.ty == b.ty
&& a.generated.as_ref().map(|value| &value.kind)
== b.generated.as_ref().map(|value| &value.kind)
});
if table.object_id() != analysis.object_id
|| engine.statistics_publication_identity(name)? != analysis.publication
|| !current.same_analysis(&analysis.maintenance)
|| !same_columns
{
return Ok(false);
}
let data_changed = engine.sampled_data_generations(name)? != analysis.data_generations;
let mut published = analysis.maintenance.clone();
published.analyzed(
table.object_id(),
analysis.row_count,
crate::statistics::value_size::FORMAT_VERSION,
analysis.sampled_at,
)?;
let Some(published) =
published.merge_sample(&analysis.maintenance, ¤t, data_changed, now_ms())?
else {
return Ok(false);
};
Self::persist_column_stats_rows(catalog, name, &analysis.statistics)?;
published.save(catalog, name)?;
*table.column_stats.write() = analysis.statistics;
table.column_stats_loaded.store(true, Ordering::Release);
table
.column_stats_dirty
.store(published.dirty(), Ordering::Release);
engine.note_table_data_changed();
Ok(true)
})
}
}