use alloc::borrow::Cow;
use alloc::string::{String, ToString};
use alloc::vec::Vec;
use spg_sql::ast::{
ColumnName, Expr, FromClause, SelectItem, SelectStatement, Statement, TableRef, UnionKind,
};
use spg_storage::{
Catalog, ColumnSchema, DataType, Row, StorageError, TableSchema, Value, VecEncoding,
};
use crate::describe;
use crate::eval::{EvalContext, EvalError};
use crate::join::RowRef;
use crate::system_catalog::collect_view_refs;
use crate::{
ByteBudget, CancelToken, Engine, EngineError, OrderKey, QueryResult, aggregate,
apply_offset_and_limit, apply_offset_and_limit_tagged, approx_row_bytes, build_order_keys,
collect_meta_view_names, collect_qualified_refs, collect_scalar_subqueries,
collect_window_nodes, compute_window_partition, eval, expr_tree_has_subquery,
materialise_in_order, materialise_meta_view, memoize, order_by_value_cmp_in, partition_key_cmp,
rewrite_window_to_columns, select_has_window, select_references_meta_view, select_refers_to,
sort_by_keys, synth_info_key_column_usage, synth_info_referential_constraints,
synth_info_routines, synth_info_statistics, synth_information_schema_columns,
synth_information_schema_tables, synth_mysql_db, synth_mysql_user, synth_pg_attribute,
synth_pg_class, synth_pg_constraint, synth_pg_database, synth_pg_extension, synth_pg_index_raw,
synth_pg_indexes, synth_pg_namespace, synth_pg_operator, synth_pg_proc, synth_pg_roles,
synth_pg_sequence, synth_pg_settings, synth_pg_timezone_abbrevs, synth_pg_timezone_names,
synth_pg_trigger, synth_pg_type, synth_pg_views, topk_trim, try_gin_jsonb_seek, try_gin_seek,
try_index_seek, try_nsw_knn, try_pk_walk_top_n, try_trgm_seek, value_is_bigint,
value_is_integer, value_to_i64,
};
struct RecursiveTermPlan<'t> {
items: Vec<&'t Expr>,
where_: Option<&'t Expr>,
alias: String,
}
fn plan_recursive_term<'t>(
t: &'t SelectStatement,
cte_name: &str,
ncols: usize,
) -> Option<RecursiveTermPlan<'t>> {
if !t.unions.is_empty()
|| !t.ctes.is_empty()
|| t.distinct
|| !t.distinct_on.is_empty()
|| t.group_by.is_some()
|| t.group_by_all
|| t.having.is_some()
|| !t.order_by.is_empty()
|| t.limit.is_some()
|| t.offset.is_some()
|| t.limit_with_ties
|| t.locking.is_some()
{
return None;
}
let from = t.from.as_ref()?;
if !from.joins.is_empty() {
return None;
}
let p = &from.primary;
if !p.name.eq_ignore_ascii_case(cte_name)
|| p.as_of_segment.is_some()
|| p.unnest_expr.is_some()
|| !p.unnest_column_aliases.is_empty()
|| p.with_ordinality
|| p.generate_series_args.is_some()
|| p.lateral_subquery.is_some()
|| p.jsonb_each_text_arg.is_some()
|| p.table_fn_call.is_some()
{
return None;
}
let unsupported = |e: &Expr| {
crate::aggregate::contains_aggregate(e)
|| crate::subquery::expr_has_subquery(e)
|| crate::window::expr_has_window_pub(e)
};
let mut items: Vec<&Expr> = Vec::with_capacity(t.items.len());
for it in &t.items {
match it {
SelectItem::Expr { expr, .. } => {
if unsupported(expr) {
return None;
}
items.push(expr);
}
_ => return None,
}
}
if items.len() != ncols {
return None;
}
if let Some(w) = &t.where_
&& unsupported(w)
{
return None;
}
Some(RecursiveTermPlan {
items,
where_: t.where_.as_ref(),
alias: p.alias.clone().unwrap_or_else(|| p.name.clone()),
})
}
impl Engine {
#[allow(
clippy::too_many_lines,
clippy::type_complexity,
clippy::needless_range_loop
)] pub(crate) fn exec_select_with_window(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
let from = stmt.from.as_ref().ok_or_else(|| {
EngineError::Unsupported("window functions require a FROM clause".into())
})?;
let (schema_cols_owned, alias_opt): (Vec<ColumnSchema>, Option<&str>);
let mut owned_rows: Vec<Row<'static>> = Vec::new();
let mut filtered: Vec<&Row<'static>> = Vec::new();
let mut rows_are_owned = false;
if from.joins.is_empty() {
let primary = &from.primary;
let is_derived = primary.lateral_subquery.is_some()
|| primary.unnest_expr.is_some()
|| primary.generate_series_args.is_some()
|| primary.jsonb_each_text_arg.is_some()
|| primary.table_fn_call.is_some();
if is_derived {
let (drows, dcols) = self.materialise_table_ref(primary)?;
schema_cols_owned = dcols;
alias_opt = primary.alias.as_deref();
let ctx = self.ev_ctx(&schema_cols_owned, alias_opt);
let mut owned: Vec<Row<'static>> = Vec::new();
for (i, row) in drows.into_iter().enumerate() {
if i.is_multiple_of(256) {
cancel.check()?;
}
if let Some(w) = &stmt.where_ {
let cond = eval::eval_expr(w, &row, &ctx)?;
if !crate::eval::predicate_is_true(&cond, "WHERE", ctx.mysql_dialect)? {
continue;
}
}
owned.push(row);
}
owned_rows = owned;
rows_are_owned = true;
} else {
let table = self.active_catalog().get(&primary.name).ok_or_else(|| {
StorageError::TableNotFound {
name: primary.name.clone(),
}
})?;
let alias = primary.alias.as_deref().unwrap_or(primary.name.as_str());
schema_cols_owned = table.schema().columns.clone();
alias_opt = Some(alias);
let ctx = self.ev_ctx(&schema_cols_owned, alias_opt);
let passes = |row: &Row<'static>| -> Result<bool, EngineError> {
if let Some(w) = &stmt.where_ {
let cond = eval::eval_expr(w, row, &ctx)?;
if !crate::eval::predicate_is_true(&cond, "WHERE", ctx.mysql_dialect)? {
return Ok(false);
}
}
Ok(true)
};
let snap = self.current_snapshot();
if table.has_cold_rows_fast() {
let mut owned: Vec<Row<'static>> = Vec::new();
for (i, row) in table.scan_visible(&snap) {
if i.is_multiple_of(256) {
cancel.check()?;
}
if passes(row)? {
owned.push(row.clone());
}
}
let hot_len = table.row_count();
for (offset, row) in self.iter_cold_rows_of_table(table).iter().enumerate() {
let i = hot_len + offset;
if i.is_multiple_of(256) {
cancel.check()?;
}
if passes(row)? {
owned.push(row.clone());
}
}
owned_rows = owned;
rows_are_owned = true;
} else {
let seek_positions: Option<Vec<usize>> = stmt.where_.as_ref().and_then(|w| {
crate::index_access::try_index_seek_positions(
w,
&schema_cols_owned,
table,
alias,
&snap,
)
});
match seek_positions {
Some(mut positions) => {
positions.sort_unstable();
for (n, pos) in positions.into_iter().enumerate() {
if n.is_multiple_of(256) {
cancel.check()?;
}
let Some(row) = table.rows().get(pos) else {
continue;
};
if passes(row)? {
filtered.push(row);
}
}
}
None => {
for (i, row) in table.scan_visible(&snap) {
if i.is_multiple_of(256) {
cancel.check()?;
}
if passes(row)? {
filtered.push(row);
}
}
}
}
}
}
} else {
let deferred = self.build_joined_filtered_rows(
from,
stmt.where_.as_ref(),
cancel,
None,
&mut ByteBudget::new(self.max_query_bytes),
)?;
owned_rows = deferred.materialise();
rows_are_owned = true;
schema_cols_owned = deferred.combined_schema;
alias_opt = None;
}
if rows_are_owned {
filtered = owned_rows.iter().collect();
}
let schema_cols = &schema_cols_owned;
let ctx = self.ev_ctx(schema_cols, alias_opt);
let alias = alias_opt.unwrap_or("");
let n_rows = filtered.len();
let mut window_nodes: Vec<Expr> = Vec::new();
for item in &stmt.items {
if let SelectItem::Expr { expr, .. } = item {
collect_window_nodes(expr, &mut window_nodes);
}
}
for o in &stmt.order_by {
collect_window_nodes(&o.expr, &mut window_nodes);
}
let mut win_vals: Vec<Vec<Value<'static>>> = Vec::with_capacity(window_nodes.len());
for wnode in &window_nodes {
let Expr::WindowFunction {
name,
args,
partition_by,
order_by,
frame,
null_treatment,
filter,
} = wnode
else {
unreachable!("collect_window_nodes pushes only WindowFunction");
};
let p_bound: Vec<Option<usize>> = partition_by
.iter()
.map(|e| crate::orderby::bound_column_position(e, schema_cols, alias_opt))
.collect();
let o_bound: Vec<Option<usize>> = order_by
.iter()
.map(|(e, _, _)| crate::orderby::bound_column_position(e, schema_cols, alias_opt))
.collect();
let arg_bound = args
.first()
.and_then(|a| crate::orderby::bound_column_position(a, schema_cols, alias_opt));
let o_colls: Vec<Option<alloc::string::String>> = o_bound
.iter()
.map(|p| {
p.and_then(|pos| schema_cols.get(pos))
.and_then(|sc| sc.collation_name.clone())
.filter(|n| crate::collate::is_supported(n))
})
.collect();
let mut indexed: Vec<(Vec<Value<'static>>, Vec<(Value, bool, Option<bool>)>, usize)> =
Vec::with_capacity(n_rows);
let int_pkey_fast = order_by.is_empty()
&& partition_by.len() == 1
&& p_bound[0].is_some_and(|pos| {
matches!(
schema_cols.get(pos).map(|c| c.ty),
Some(
spg_storage::DataType::Int
| spg_storage::DataType::BigInt
| spg_storage::DataType::SmallInt
)
)
});
let int_okey_fast = partition_by.is_empty()
&& order_by.len() == 1
&& frame.is_none()
&& filter.is_none()
&& matches!(null_treatment, spg_sql::ast::NullTreatment::Respect)
&& name.eq_ignore_ascii_case("row_number")
&& o_bound[0].is_some_and(|pos| {
matches!(
schema_cols.get(pos).map(|c| c.ty),
Some(
spg_storage::DataType::Int
| spg_storage::DataType::BigInt
| spg_storage::DataType::SmallInt
)
)
});
let mut int_okey_bailed = false;
if int_okey_fast {
let pos = o_bound[0].expect("gated bound");
let desc = order_by[0].1;
let nulls_first = order_by[0].2.unwrap_or(desc);
let mut keyed: Vec<(bool, i64, usize)> = Vec::with_capacity(n_rows);
for (i, row) in filtered.iter().enumerate() {
match row.values.get(pos) {
Some(Value::Int(n)) => keyed.push((false, i64::from(*n), i)),
Some(Value::BigInt(n)) => keyed.push((false, *n, i)),
Some(Value::SmallInt(n)) => keyed.push((false, i64::from(*n), i)),
Some(Value::Null) | None => keyed.push((true, 0, i)),
Some(_) => {
int_okey_bailed = true;
break;
}
}
}
if !int_okey_bailed {
let null_rank = |is_null: bool| -> u8 { u8::from(is_null != nulls_first) };
keyed.sort_unstable_by(|a, b| {
null_rank(a.0)
.cmp(&null_rank(b.0))
.then_with(|| {
if a.0 {
core::cmp::Ordering::Equal
} else if desc {
b.1.cmp(&a.1)
} else {
a.1.cmp(&b.1)
}
})
.then_with(|| a.2.cmp(&b.2))
});
for (_, _, i) in keyed {
indexed.push((Vec::new(), Vec::new(), i));
}
} else {
indexed.clear();
}
}
if int_okey_fast && !int_okey_bailed {
} else if int_pkey_fast {
let pos = p_bound[0].expect("gated bound");
let mut slot: hashbrown::HashMap<Option<i64>, usize> = hashbrown::HashMap::new();
let mut groups: Vec<Vec<usize>> = Vec::new();
for (i, row) in filtered.iter().enumerate() {
let k: Option<i64> = match row.values.get(pos) {
Some(Value::BigInt(n)) => Some(*n),
Some(Value::Int(n)) => Some(i64::from(*n)),
Some(Value::SmallInt(n)) => Some(i64::from(*n)),
_ => None,
};
match slot.get(&k) {
Some(&gi) => groups[gi].push(i),
None => {
slot.insert(k, groups.len());
groups.push(alloc::vec![i]);
}
}
}
for g in groups {
for i in g {
let k: Value<'static> = match filtered[i].values.get(pos) {
Some(v) => v.clone(),
None => Value::Null,
};
indexed.push((alloc::vec![k], Vec::new(), i));
}
}
} else {
for (i, row) in filtered.iter().enumerate() {
let pkey: Vec<Value<'static>> = partition_by
.iter()
.enumerate()
.map(
|(k, p)| match p_bound[k].and_then(|pos| row.values.get(pos)) {
Some(v) => Ok(v.clone()),
None => eval::eval_expr(p, row, &ctx),
},
)
.collect::<Result<_, _>>()?;
let okey: Vec<(Value, bool, Option<bool>)> = order_by
.iter()
.enumerate()
.map(|(k, (e, desc, nf))| -> Result<_, EngineError> {
let v = match o_bound[k].and_then(|pos| row.values.get(pos)) {
Some(v) => v.clone(),
None => eval::eval_expr(e, row, &ctx)?,
};
let v = match crate::orderby::enum_order_ordinal(e, &v, &ctx) {
Some(ord) => Value::Float(ord),
None => v,
};
Ok((v, *desc, *nf))
})
.collect::<Result<_, _>>()?;
indexed.push((pkey, okey, i));
}
}
if int_okey_fast && !int_okey_bailed {
} else if int_pkey_fast {
} else if order_by.is_empty() && !partition_by.is_empty() {
let mut slot: hashbrown::HashMap<String, usize> = hashbrown::HashMap::new();
let mut groups: Vec<
Vec<(Vec<Value<'static>>, Vec<(Value, bool, Option<bool>)>, usize)>,
> = Vec::new();
let mut keybuf = String::new();
for entry in indexed.drain(..) {
keybuf.clear();
for v in &entry.0 {
crate::aggregate::push_canonical_key(&mut keybuf, v);
}
match slot.get(keybuf.as_str()) {
Some(&gi) => groups[gi].push(entry),
None => {
slot.insert(keybuf.clone(), groups.len());
groups.push(alloc::vec![entry]);
}
}
}
for g in groups {
indexed.extend(g);
}
} else {
indexed.sort_by(|a, b| {
let p_cmp = partition_key_cmp(&a.0, &b.0);
if p_cmp != core::cmp::Ordering::Equal {
return p_cmp;
}
crate::window::order_key_cmp_in(&a.1, &b.1, &o_colls)
});
}
let mut out_vals: Vec<Value<'static>> = alloc::vec![Value::Null; n_rows];
let mut p_start = 0;
while p_start < indexed.len() {
let mut p_end = p_start + 1;
while p_end < indexed.len()
&& partition_key_cmp(&indexed[p_start].0, &indexed[p_end].0)
== core::cmp::Ordering::Equal
{
p_end += 1;
}
compute_window_partition(
name,
args,
arg_bound,
!order_by.is_empty(),
frame.as_ref(),
*null_treatment,
filter.as_deref(),
&indexed[p_start..p_end],
&filtered,
&ctx,
&mut out_vals,
)?;
p_start = p_end;
}
win_vals.push(out_vals);
}
let mut ext_cols = schema_cols.clone();
for i in 0..window_nodes.len() {
ext_cols.push(ColumnSchema::new(
alloc::format!("__win_{i}"),
DataType::Text, true,
));
}
let mut rewritten_items: Vec<SelectItem> = Vec::with_capacity(stmt.items.len());
for item in &stmt.items {
let new_item = match item {
SelectItem::Wildcard => SelectItem::Wildcard,
SelectItem::QualifiedWildcard(q) => SelectItem::QualifiedWildcard(q.clone()),
SelectItem::Expr { expr, alias } => {
let mut e = expr.clone();
rewrite_window_to_columns(&mut e, &window_nodes);
let alias = if alias.is_none() && e != *expr {
Some(default_output_name(expr, self.backslash_escapes))
} else {
alias.clone()
};
SelectItem::Expr { expr: e, alias }
}
};
rewritten_items.push(new_item);
}
let ext_ctx = self.ev_ctx(&ext_cols, alias_opt);
let projection = build_projection_hiding_tail(
&rewritten_items,
&ext_cols,
alias,
self.backslash_escapes,
window_nodes.len(),
)?;
let mut tagged: Vec<(Vec<OrderKey>, Row)> = Vec::with_capacity(n_rows);
let mut ext_row: Row<'static> =
Row::new(Vec::with_capacity(schema_cols.len() + window_nodes.len()));
for i in 0..n_rows {
if i.is_multiple_of(256) {
cancel.check()?;
}
ext_row.values.clear();
ext_row.values.extend(filtered[i].values.iter().cloned());
for w in 0..window_nodes.len() {
ext_row.values.push(win_vals[w][i].clone());
}
let row = &ext_row;
let mut values = Vec::with_capacity(projection.len());
for p in &projection {
values.push(eval::eval_expr(&p.expr, row, &ext_ctx)?);
}
let order_keys = if stmt.order_by.is_empty() {
Vec::new()
} else {
let mut keys = Vec::with_capacity(stmt.order_by.len());
for o in &stmt.order_by {
let mut e = o.expr.clone();
rewrite_window_to_columns(&mut e, &window_nodes);
let key = eval::eval_expr(&e, row, &ext_ctx)?;
match crate::orderby::enum_order_ordinal(&e, &key, &ext_ctx) {
Some(ord) => keys.push(value_to_order_key(&Value::Float(ord))?),
None => keys.push(value_to_order_key(&key)?),
}
}
keys
};
tagged.push((order_keys, Row::new(values)));
}
if !stmt.order_by.is_empty() {
let descs: Vec<bool> = stmt.order_by.iter().map(|o| o.desc).collect();
sort_by_keys(&mut tagged, &descs);
}
let mut out_rows: Vec<Row<'static>> = tagged.into_iter().map(|(_, r)| r).collect();
if stmt.distinct {
out_rows = dedup_rows(out_rows, self.backslash_escapes);
}
apply_offset_and_limit(&mut out_rows, stmt.offset_literal(), stmt.limit_literal());
let final_cols: Vec<ColumnSchema> = projection
.into_iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name, p.ty, p.nullable);
c.user_enum_type = p.user_enum_type;
c.collation_name = p.collation_name;
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
Ok(QueryResult::Rows {
columns: final_cols,
rows: out_rows,
})
}
pub(crate) fn exec_select_with_meta_views(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
let catalog = self.meta_view_catalog(stmt)?;
let mut temp = Engine::restore(catalog);
if let Some(c) = self.clock {
temp = temp.with_clock(c);
}
if let Some(f) = self.salt_fn {
temp = temp.with_salt_fn(f);
}
temp.session_params.clone_from(&self.session_params);
temp.users.clone_from(&self.users);
temp.backslash_escapes = self.backslash_escapes;
temp.mysql_strict = self.mysql_strict;
temp.render_style = self.render_style;
temp.tz_offset_fn = self.tz_offset_fn;
temp.tz_localize_fn = self.tz_localize_fn;
temp.tz_abbrev_fn = self.tz_abbrev_fn;
temp.meta_views_materialised = true;
temp.exec_select_cancel(stmt, cancel)
}
pub(crate) fn meta_view_catalog(&self, stmt: &SelectStatement) -> Result<Catalog, EngineError> {
let mut needed: alloc::collections::BTreeSet<String> = alloc::collections::BTreeSet::new();
collect_meta_view_names(stmt, &mut needed);
let mut catalog = self.active_catalog().clone();
for view in &needed {
if catalog.get(view).is_some() {
continue;
}
match view.as_str() {
"__spg_info_columns" => {
let (schema, rows) = synth_information_schema_columns(
self.active_catalog(),
self.backslash_escapes,
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_tables" => {
let (schema, rows) = synth_information_schema_tables(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_class" => {
let (schema, rows) = synth_pg_class(
self.active_catalog(),
i64::try_from(self.vacuum_oldest_active()).unwrap_or(i64::MAX),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_attribute" => {
let (schema, rows) = synth_pg_attribute(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_type" => {
let (schema, rows) = synth_pg_type(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_operator" => {
let (schema, rows) = synth_pg_operator(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_proc" => {
let (schema, rows) = synth_pg_proc(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_trigger" => {
let (schema, rows) = synth_pg_trigger(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_namespace" => {
let (schema, rows) = synth_pg_namespace(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_tables" => {
let (schema, rows) =
crate::system_catalog::synth_pg_tables(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_enum" => {
let (schema, rows) =
crate::system_catalog::synth_pg_enum(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_prepared_statements" => {
let (schema, rows) = crate::system_catalog::synth_pg_prepared_statements(
&self.prepared_statements,
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_replication_slots" => {
let (schema, rows) =
crate::system_catalog::synth_pg_replication_slots(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_publication" => {
let (schema, rows) = crate::system_catalog::synth_pg_publication(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_subscription" => {
let (schema, rows) = crate::system_catalog::synth_pg_subscription(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_database" => {
let (schema, rows) = crate::system_catalog::synth_pg_stat_database(
self,
self.stat_tup_inserted,
self.stat_tup_updated,
self.stat_tup_deleted,
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_user_tables" => {
let (schema, rows) = crate::system_catalog::synth_pg_stat_user_tables(
self.active_catalog(),
&self.table_write_stats,
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_user_indexes" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_user_indexes(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_bgwriter" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_bgwriter(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_checkpointer" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_checkpointer(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_wal" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_wal(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_slru" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_slru(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_subscription_stats" => {
let (schema, rows) = crate::system_catalog::synth_pg_stat_subscription_stats(
self.active_catalog(),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_archiver" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_archiver(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_replication" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_replication(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_am" => {
let (schema, rows) = crate::system_catalog::synth_pg_am(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_io" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_io(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_user_functions" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_user_functions(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_largeobject" => {
let (schema, rows) =
crate::system_catalog::synth_pg_largeobject(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_largeobject_metadata" => {
let (schema, rows) =
crate::system_catalog::synth_pg_largeobject_metadata(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_statistic_ext" => {
let (schema, rows) =
crate::system_catalog::synth_pg_statistic_ext(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_statistic" => {
let (schema, rows) =
crate::system_catalog::synth_pg_statistic(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_progress_vacuum" => {
let (schema, rows) =
crate::system_catalog::synth_pg_stat_progress_vacuum(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_progress_create_index" => {
let (schema, rows) = crate::system_catalog::synth_pg_stat_progress_create_index(
self.active_catalog(),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_stat_progress_analyze" => {
let (schema, rows) = crate::system_catalog::synth_pg_stat_progress_analyze(
self.active_catalog(),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_inherits" => {
let (schema, rows) =
crate::system_catalog::synth_pg_inherits(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_ts_config_map" => {
let (schema, rows) =
crate::system_catalog::synth_pg_ts_config_map(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_ts_config" => {
let (schema, rows) =
crate::system_catalog::synth_pg_ts_config(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_ts_dict" => {
let (schema, rows) =
crate::system_catalog::synth_pg_ts_dict(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_ts_parser" => {
let (schema, rows) =
crate::system_catalog::synth_pg_ts_parser(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_ts_template" => {
let (schema, rows) =
crate::system_catalog::synth_pg_ts_template(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_depend" => {
let (schema, rows) =
crate::system_catalog::synth_pg_depend(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_attrdef" => {
let (schema, rows) =
crate::system_catalog::synth_pg_attrdef(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_policy" => {
let (schema, rows) =
crate::system_catalog::synth_pg_policy(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_policies" => {
let (schema, rows) =
crate::system_catalog::synth_pg_policies(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_collation" => {
let (schema, rows) =
crate::system_catalog::synth_pg_collation(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_tablespace" => {
let (schema, rows) =
crate::system_catalog::synth_pg_tablespace(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_indexes" => {
let (schema, rows) = synth_pg_indexes(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_description" => {
let (schema, rows) =
crate::system_catalog::synth_pg_description(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_index" => {
let (schema, rows) = synth_pg_index_raw(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_constraint" => {
let (schema, rows) = synth_pg_constraint(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_sequence" => {
let (schema, rows) = synth_pg_sequence(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_database" => {
let (schema, rows) = synth_pg_database(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_roles" => {
let (schema, rows) = synth_pg_roles(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_user" => {
let (schema, rows) = crate::system_catalog::synth_pg_user(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_auth_members" => {
let (schema, rows) = crate::system_catalog::synth_pg_auth_members(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_views" => {
let (schema, rows) = synth_pg_views(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_rules" => {
let (schema, rows) =
crate::system_catalog::synth_pg_rules(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_rewrite" => {
let (schema, rows) =
crate::system_catalog::synth_pg_rewrite(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_matviews" => {
let (schema, rows) =
crate::system_catalog::synth_pg_matviews(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_db_role_setting" => {
let (schema, rows) = crate::system_catalog::synth_pg_db_role_setting(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_language" => {
let (schema, rows) = crate::system_catalog::synth_pg_language();
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_sequences" => {
let (schema, rows) =
crate::system_catalog::synth_pg_sequences(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_range" => {
let (schema, rows) = crate::system_catalog::synth_pg_range();
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_partitioned_table" => {
let (schema, rows) =
crate::system_catalog::synth_pg_partitioned_table(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_authid" => {
let (schema, rows) = crate::system_catalog::synth_pg_authid(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_group" => {
let (schema, rows) = crate::system_catalog::synth_pg_group(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_shadow" => {
let (schema, rows) = crate::system_catalog::synth_pg_shadow(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_cast" => {
let (schema, rows) = crate::system_catalog::synth_pg_cast();
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_foreign_table" => {
let (schema, rows) = crate::system_catalog::synth_pg_foreign_table();
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_extension" => {
let (schema, rows) = synth_pg_extension();
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_timezone_names" => {
let (schema, rows) = synth_pg_timezone_names(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_timezone_abbrevs" => {
let (schema, rows) = synth_pg_timezone_abbrevs(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_pg_settings" => {
let (schema, rows) = synth_pg_settings(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_column_privileges" => {
let (schema, rows) =
crate::system_catalog::synth_info_column_privileges(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_role_table_grants" | "__spg_info_table_privileges" => {
let grantee = self.current_role().to_string();
let (schema, rows) = crate::system_catalog::synth_info_role_table_grants(
self.active_catalog(),
&grantee,
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_key_column_usage" => {
let (schema, rows) = synth_info_key_column_usage(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_referential_constraints" => {
let (schema, rows) = synth_info_referential_constraints(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_statistics" => {
let (schema, rows) = synth_info_statistics(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_routines" => {
let (schema, rows) = synth_info_routines();
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_attributes" => {
let (schema, rows) = crate::system_catalog::synth_information_schema_attributes(
self.active_catalog(),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_domains" => {
let (schema, rows) = crate::system_catalog::synth_information_schema_domains(
self.active_catalog(),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_schemata" => {
let (schema, rows) = crate::system_catalog::synth_information_schema_schemata(
self.active_catalog(),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_views" => {
let (schema, rows) = crate::system_catalog::synth_information_schema_views(
self.active_catalog(),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_table_constraints" => {
let (schema, rows) =
crate::system_catalog::synth_information_schema_table_constraints(
self.active_catalog(),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_constraint_column_usage" => {
let (schema, rows) = crate::system_catalog::synth_info_constraint_column_usage(
self.active_catalog(),
);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_triggers" => {
let (schema, rows) =
crate::system_catalog::synth_info_triggers(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_check_constraints" => {
let (schema, rows) =
crate::system_catalog::synth_info_check_constraints(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_info_sequences" => {
let (schema, rows) =
crate::system_catalog::synth_info_sequences(self.active_catalog());
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_mysql_user" => {
let (schema, rows) = synth_mysql_user(self);
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
"__spg_mysql_db" => {
let (schema, rows) = synth_mysql_db();
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
other if crate::system_catalog::synth_empty_pg_catalog(other).is_some() => {
let (schema, rows) =
crate::system_catalog::synth_empty_pg_catalog(other).expect("just checked");
materialise_meta_view(&mut catalog, view, schema, rows)?;
}
_ => {
return Err(EngineError::Unsupported(alloc::format!(
"meta view {view:?} is not yet materialisable; \
v7.16.2 covers information_schema.columns / .tables \
and pg_catalog.pg_class / pg_attribute; \
v7.17.0 P0-50..P0-57 add pg_type / pg_proc / pg_namespace / \
pg_indexes / pg_index / pg_constraint / pg_database / pg_roles / \
pg_user / pg_views / pg_matviews / pg_settings"
)));
}
}
}
Ok(catalog)
}
pub(crate) fn exec_with_ctes(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
cancel.check()?;
if stmt.ctes.iter().any(|c| c.body.is_modifying()) {
return Err(EngineError::Unsupported(
"WITH clause containing a data-modifying statement must be at the top level".into(),
));
}
let catalog = self.materialise_ctes_readonly(&stmt.ctes, cancel)?;
let mut body = stmt.clone();
body.ctes = Vec::new();
let mut temp = Engine::restore(catalog);
if let Some(c) = self.clock {
temp = temp.with_clock(c);
}
if let Some(f) = self.salt_fn {
temp = temp.with_salt_fn(f);
}
temp.exec_select_cancel(&body, cancel)
}
pub(crate) fn materialise_ctes_readonly(
&self,
ctes: &[spg_sql::ast::Cte],
cancel: CancelToken<'_>,
) -> Result<crate::Catalog, EngineError> {
cancel.check()?;
let mut catalog = self.active_catalog().clone();
for cte in ctes {
let body_select = cte.body.as_select().ok_or_else(|| {
EngineError::Unsupported(alloc::format!(
"data-modifying CTE not supported on this SELECT entry"
))
})?;
let (columns, rows) = if cte.recursive && select_refers_to(body_select, &cte.name) {
let synthetic = spg_sql::ast::Cte {
name: cte.name.clone(),
body: spg_sql::ast::CteBody::Select(body_select.clone()),
recursive: true,
column_overrides: cte.column_overrides.clone(),
search: None,
cycle: None,
};
if catalog.get(&cte.name).is_some() {
let _ = catalog.drop_table(&cte.name);
}
self.materialise_recursive_cte(&synthetic, &catalog, cancel)?
} else {
let mut cte_engine = Engine::restore(catalog.clone());
if let Some(c) = self.clock {
cte_engine = cte_engine.with_clock(c);
}
if let Some(f) = self.salt_fn {
cte_engine = cte_engine.with_salt_fn(f);
}
let body_result = cte_engine.exec_select_cancel(body_select, cancel)?;
let QueryResult::Rows { columns, rows } = body_result else {
return Err(EngineError::Unsupported(alloc::format!(
"CTE {:?} body did not return rows",
cte.name
)));
};
(columns, rows)
};
let inferred = infer_column_types(&columns, &rows);
let mut columns = inferred;
if !cte.column_overrides.is_empty() {
if cte.column_overrides.len() != columns.len() {
return Err(EngineError::Unsupported(alloc::format!(
"CTE {:?} column list has {} names but body returns {} columns",
cte.name,
cte.column_overrides.len(),
columns.len()
)));
}
for (col, name) in columns.iter_mut().zip(cte.column_overrides.iter()) {
col.name.clone_from(name);
}
}
let schema = TableSchema::new(cte.name.clone(), columns);
if catalog.get(&cte.name).is_some() {
let _ = catalog.drop_table(&cte.name);
}
catalog.create_table(schema).map_err(EngineError::Storage)?;
let table = catalog
.get_mut(&cte.name)
.expect("just-created CTE table must exist");
for row in rows {
table.insert(row).map_err(EngineError::Storage)?;
}
}
Ok(catalog)
}
#[allow(dead_code)]
pub(crate) fn materialise_ctes(
&mut self,
ctes: &[spg_sql::ast::Cte],
cancel: CancelToken<'_>,
) -> Result<crate::Catalog, EngineError> {
cancel.check()?;
let mut catalog = self.active_catalog().clone();
for cte in ctes {
let body_target = match &cte.body {
spg_sql::ast::CteBody::Select(_) => None,
spg_sql::ast::CteBody::Insert(i) => Some(i.table.as_str()),
spg_sql::ast::CteBody::Update(u) => Some(u.table.as_str()),
spg_sql::ast::CteBody::Delete(d) => Some(d.table.as_str()),
spg_sql::ast::CteBody::Merge(m) => Some(m.target.as_str()),
};
if let Some(t) = body_target
&& ctes.iter().any(|c| c.name.eq_ignore_ascii_case(t))
&& catalog.get(t).is_none()
{
return Err(EngineError::Storage(
spg_storage::StorageError::TableNotFound { name: t.into() },
));
}
}
for cte in ctes {
if catalog.get(&cte.name).is_some() {
return Err(EngineError::Unsupported(alloc::format!(
"CTE name {:?} shadows an existing table; rename the CTE",
cte.name
)));
}
let (columns, rows) = match &cte.body {
spg_sql::ast::CteBody::Select(body)
if cte.recursive && select_refers_to(body, &cte.name) =>
{
let synthetic = spg_sql::ast::Cte {
name: cte.name.clone(),
body: spg_sql::ast::CteBody::Select(body.clone()),
recursive: true,
column_overrides: cte.column_overrides.clone(),
search: None,
cycle: None,
};
self.materialise_recursive_cte(&synthetic, &catalog, cancel)?
}
spg_sql::ast::CteBody::Select(body) => {
let mut cte_engine = Engine::restore(catalog.clone());
if let Some(c) = self.clock {
cte_engine = cte_engine.with_clock(c);
}
if let Some(f) = self.salt_fn {
cte_engine = cte_engine.with_salt_fn(f);
}
let body_result = cte_engine.exec_select_cancel(body, cancel)?;
let QueryResult::Rows { columns, rows } = body_result else {
return Err(EngineError::Unsupported(alloc::format!(
"CTE {:?} body did not return rows",
cte.name
)));
};
(columns, rows)
}
spg_sql::ast::CteBody::Insert(body) => {
self.exec_modifying_cte_insert(&cte.name, body, cancel)?
}
spg_sql::ast::CteBody::Update(body) => {
self.exec_modifying_cte_update(&cte.name, body, cancel)?
}
spg_sql::ast::CteBody::Delete(body) => {
self.exec_modifying_cte_delete(&cte.name, body, cancel)?
}
spg_sql::ast::CteBody::Merge(body) => {
self.exec_modifying_cte_merge(&cte.name, body, cancel)?
}
};
let inferred = infer_column_types(&columns, &rows);
let mut columns = inferred;
if !cte.column_overrides.is_empty() {
if cte.column_overrides.len() != columns.len() {
return Err(EngineError::Unsupported(alloc::format!(
"CTE {:?} column list has {} names but body returns {} columns",
cte.name,
cte.column_overrides.len(),
columns.len()
)));
}
for (col, name) in columns.iter_mut().zip(cte.column_overrides.iter()) {
col.name.clone_from(name);
}
}
let schema = TableSchema::new(cte.name.clone(), columns);
catalog.create_table(schema).map_err(EngineError::Storage)?;
let table = catalog
.get_mut(&cte.name)
.expect("just-created CTE table must exist");
for row in rows {
table.insert(row).map_err(EngineError::Storage)?;
}
}
Ok(catalog)
}
fn exec_modifying_cte_insert(
&mut self,
cte_name: &str,
body: &spg_sql::ast::InsertStatement,
_cancel: CancelToken<'_>,
) -> Result<
(
Vec<spg_storage::ColumnSchema>,
Vec<spg_storage::Row<'static>>,
),
EngineError,
> {
let body = body.clone();
let result = self.exec_insert(body)?;
match result {
QueryResult::Rows { columns, rows } => Ok((columns, rows)),
QueryResult::CommandOk { .. } => {
let placeholder = spg_storage::ColumnSchema::new(
alloc::format!("{cte_name}_returning_absent"),
spg_storage::DataType::Text,
true,
);
Ok((alloc::vec![placeholder], Vec::new()))
}
}
}
fn exec_modifying_cte_update(
&mut self,
cte_name: &str,
body: &spg_sql::ast::UpdateStatement,
cancel: CancelToken<'_>,
) -> Result<
(
Vec<spg_storage::ColumnSchema>,
Vec<spg_storage::Row<'static>>,
),
EngineError,
> {
let body = body.clone();
let result = self.exec_update_cancel(&body, cancel)?;
match result {
QueryResult::Rows { columns, rows } => Ok((columns, rows)),
QueryResult::CommandOk { .. } => {
let placeholder = spg_storage::ColumnSchema::new(
alloc::format!("{cte_name}_returning_absent"),
spg_storage::DataType::Text,
true,
);
Ok((alloc::vec![placeholder], Vec::new()))
}
}
}
fn exec_modifying_cte_delete(
&mut self,
cte_name: &str,
body: &spg_sql::ast::DeleteStatement,
cancel: CancelToken<'_>,
) -> Result<
(
Vec<spg_storage::ColumnSchema>,
Vec<spg_storage::Row<'static>>,
),
EngineError,
> {
let body = body.clone();
let result = self.exec_delete_cancel(&body, cancel)?;
match result {
QueryResult::Rows { columns, rows } => Ok((columns, rows)),
QueryResult::CommandOk { .. } => {
let placeholder = spg_storage::ColumnSchema::new(
alloc::format!("{cte_name}_returning_absent"),
spg_storage::DataType::Text,
true,
);
Ok((alloc::vec![placeholder], Vec::new()))
}
}
}
fn exec_modifying_cte_merge(
&mut self,
cte_name: &str,
body: &spg_sql::ast::MergeStatement,
cancel: CancelToken<'_>,
) -> Result<
(
Vec<spg_storage::ColumnSchema>,
Vec<spg_storage::Row<'static>>,
),
EngineError,
> {
let body = body.clone();
let result = self.exec_merge_cancel(&body, cancel)?;
match result {
QueryResult::Rows { columns, rows } => Ok((columns, rows)),
QueryResult::CommandOk { .. } => {
let placeholder = spg_storage::ColumnSchema::new(
alloc::format!("{cte_name}_returning_absent"),
spg_storage::DataType::Text,
true,
);
Ok((alloc::vec![placeholder], Vec::new()))
}
}
}
#[allow(clippy::too_many_lines)]
pub(crate) fn materialise_recursive_cte(
&self,
cte: &spg_sql::ast::Cte,
base_catalog: &Catalog,
cancel: CancelToken<'_>,
) -> Result<(Vec<ColumnSchema>, Vec<Row<'static>>), EngineError> {
const MAX_TOTAL_ROWS: usize = 1_000_000;
const MAX_ITERATIONS: usize = 100_000;
cancel.check()?;
let body_select = cte.body.as_select().ok_or_else(|| {
EngineError::Unsupported(alloc::format!(
"WITH RECURSIVE {:?} body must be a SELECT, not a data-modifying statement",
cte.name
))
})?;
if body_select.unions.is_empty() {
return Err(EngineError::Unsupported(alloc::format!(
"WITH RECURSIVE {:?} body must be a UNION of an anchor and a recursive term",
cte.name
)));
}
let mut anchor = body_select.clone();
let all_union_terms = core::mem::take(&mut anchor.unions);
anchor.ctes = Vec::new();
let (anchor_terms, union_terms): (Vec<_>, Vec<_>) = all_union_terms
.into_iter()
.partition(|(_, t)| !select_refers_to(t, &cte.name));
let anchor_result = self.exec_select_cancel(&anchor, cancel)?;
let QueryResult::Rows {
columns: anchor_cols,
rows: mut anchor_rows,
} = anchor_result
else {
return Err(EngineError::Unsupported(alloc::format!(
"WITH RECURSIVE {:?}: anchor did not return rows",
cte.name
)));
};
for (_, term) in &anchor_terms {
let mut term = term.clone();
term.ctes = Vec::new();
if let QueryResult::Rows { rows, .. } = self.exec_select_cancel(&term, cancel)? {
anchor_rows.extend(rows);
}
}
let mut columns = infer_column_types(&anchor_cols, &anchor_rows);
if !cte.column_overrides.is_empty() {
if cte.column_overrides.len() != columns.len() {
return Err(EngineError::Unsupported(alloc::format!(
"CTE {:?} column list has {} names but anchor returns {} columns",
cte.name,
cte.column_overrides.len(),
columns.len()
)));
}
for (col, name) in columns.iter_mut().zip(cte.column_overrides.iter()) {
col.name.clone_from(name);
}
}
let mut all_rows: Vec<Row<'static>> = anchor_rows.clone();
let mut working_set: Vec<Row<'static>> = anchor_rows;
let mut seen: alloc::collections::BTreeSet<Vec<u8>> = alloc::collections::BTreeSet::new();
let all_union_all = union_terms.iter().all(|(k, _)| matches!(k, UnionKind::All));
if !all_union_all {
for r in &all_rows {
seen.insert(encode_row_key(r));
}
}
let mut iter_catalog = base_catalog.clone();
let schema = TableSchema::new(cte.name.clone(), columns.clone());
iter_catalog
.create_table(schema)
.map_err(EngineError::Storage)?;
let mut iter_engine = Engine::restore(iter_catalog);
if let Some(c) = self.clock {
iter_engine = iter_engine.with_clock(c);
}
if let Some(f) = self.salt_fn {
iter_engine = iter_engine.with_salt_fn(f);
}
let recursive_terms: Vec<SelectStatement> = union_terms
.iter()
.map(|(_, t)| {
let mut t = t.clone();
t.ctes = Vec::new();
t
})
.collect();
let term_plans: Option<Vec<RecursiveTermPlan<'_>>> = recursive_terms
.iter()
.map(|t| plan_recursive_term(t, &cte.name, columns.len()))
.collect();
let fast_ctx = term_plans.as_ref().map(|plans| {
let alias = plans[0].alias.clone();
(alias, ())
});
for iter in 0..MAX_ITERATIONS {
cancel.check()?;
if working_set.is_empty() {
break;
}
if let (Some(plans), Some((_, ()))) = (term_plans.as_ref(), fast_ctx.as_ref()) {
let mut next_set: Vec<Row<'static>> = Vec::new();
for plan in plans {
let ctx = self.ev_ctx(&columns, Some(&plan.alias));
for row in &working_set {
cancel.check()?;
if let Some(w) = plan.where_ {
let v = eval::eval_expr(w, row, &ctx).map_err(EngineError::Eval)?;
if !matches!(v, Value::Bool(true)) {
continue;
}
}
let mut vals: Vec<Value<'static>> = Vec::with_capacity(plan.items.len());
for it in &plan.items {
vals.push(eval::eval_expr(it, row, &ctx).map_err(EngineError::Eval)?);
}
let out = Row::new(vals);
if !all_union_all {
let key = encode_row_key(&out);
if !seen.insert(key) {
continue;
}
}
next_set.push(out);
}
}
if next_set.is_empty() {
break;
}
all_rows.extend(next_set.iter().cloned());
working_set = next_set;
if all_rows.len() > MAX_TOTAL_ROWS {
return Err(EngineError::Unsupported(alloc::format!(
"WITH RECURSIVE {:?}: produced more than {MAX_TOTAL_ROWS} rows — likely runaway recursion",
cte.name
)));
}
if iter + 1 == MAX_ITERATIONS {
return Err(EngineError::Unsupported(alloc::format!(
"WITH RECURSIVE {:?}: exceeded {MAX_ITERATIONS} iterations",
cte.name
)));
}
continue;
}
{
let cat = iter_engine.base_catalog_mut();
let table = cat.get_mut(&cte.name).expect("created above");
table.truncate();
for row in &working_set {
table.insert(row.clone()).map_err(EngineError::Storage)?;
}
}
let mut next_set: Vec<Row<'static>> = Vec::new();
for term in &recursive_terms {
let r = iter_engine.exec_select_cancel(term, cancel)?;
let QueryResult::Rows {
columns: rc,
rows: rs,
} = r
else {
return Err(EngineError::Unsupported(alloc::format!(
"WITH RECURSIVE {:?}: recursive term did not return rows",
cte.name
)));
};
if rc.len() != columns.len() {
return Err(EngineError::Unsupported(alloc::format!(
"WITH RECURSIVE {:?}: column count of recursive term ({}) does not match anchor ({})",
cte.name,
rc.len(),
columns.len()
)));
}
for row in rs {
if !all_union_all {
let key = encode_row_key(&row);
if !seen.insert(key) {
continue;
}
}
next_set.push(row);
}
}
if next_set.is_empty() {
break;
}
all_rows.extend(next_set.iter().cloned());
working_set = next_set;
if all_rows.len() > MAX_TOTAL_ROWS {
return Err(EngineError::Unsupported(alloc::format!(
"WITH RECURSIVE {:?}: produced more than {MAX_TOTAL_ROWS} rows — likely runaway recursion",
cte.name
)));
}
if iter + 1 == MAX_ITERATIONS {
return Err(EngineError::Unsupported(alloc::format!(
"WITH RECURSIVE {:?}: exceeded {MAX_ITERATIONS} iterations",
cte.name
)));
}
}
Ok((columns, all_rows))
}
pub(crate) fn resolve_select_subqueries(
&self,
stmt: &mut SelectStatement,
cancel: CancelToken<'_>,
) -> Result<(), EngineError> {
for item in &mut stmt.items {
if let SelectItem::Expr { expr, alias } = item {
if alias.is_none()
&& matches!(
expr,
Expr::ScalarSubquery(_)
| Expr::Exists { .. }
| Expr::InSubquery { .. }
| Expr::RowInSubquery { .. }
| Expr::RowCmpSubquery { .. }
)
{
*alias = Some(default_output_name(expr, self.backslash_escapes));
}
self.resolve_expr_subqueries(expr, cancel)?;
}
}
if let Some(w) = &mut stmt.where_ {
self.resolve_expr_subqueries(w, cancel)?;
}
if let Some(from) = &mut stmt.from {
for j in &mut from.joins {
if let Some(on) = &mut j.on {
self.resolve_expr_subqueries(on, cancel)?;
}
}
}
if let Some(gs) = &mut stmt.group_by {
for g in gs {
self.resolve_expr_subqueries(g, cancel)?;
}
}
if let Some(h) = &mut stmt.having {
self.resolve_expr_subqueries(h, cancel)?;
}
for o in &mut stmt.order_by {
self.resolve_expr_subqueries(&mut o.expr, cancel)?;
}
for (_, peer) in &mut stmt.unions {
self.resolve_select_subqueries(peer, cancel)?;
}
Ok(())
}
#[allow(clippy::only_used_in_recursion)] pub(crate) fn resolve_expr_subqueries(
&self,
e: &mut Expr,
cancel: CancelToken<'_>,
) -> Result<(), EngineError> {
if let Some(replacement) = self.subquery_replacement(e, cancel)? {
*e = replacement;
return Ok(());
}
match e {
Expr::NamedArg { expr, .. } => self.resolve_expr_subqueries(expr, cancel)?,
Expr::Variadic(expr) => self.resolve_expr_subqueries(expr, cancel)?,
Expr::AggregateOrdered { call, order_by, .. } => {
self.resolve_expr_subqueries(call, cancel)?;
for o in order_by.iter_mut() {
self.resolve_expr_subqueries(&mut o.expr, cancel)?;
}
}
Expr::Binary { lhs, rhs, .. } => {
self.resolve_expr_subqueries(lhs, cancel)?;
self.resolve_expr_subqueries(rhs, cancel)?;
}
Expr::Unary { expr, .. }
| Expr::Cast { expr, .. }
| Expr::IsNull { expr, .. }
| Expr::BoolTest { expr, .. }
| Expr::FieldAccess { base: expr, .. } => {
self.resolve_expr_subqueries(expr, cancel)?;
}
Expr::FunctionCall { args, .. } => {
for a in args {
self.resolve_expr_subqueries(a, cancel)?;
}
}
Expr::Like { expr, pattern, .. } => {
self.resolve_expr_subqueries(expr, cancel)?;
self.resolve_expr_subqueries(pattern, cancel)?;
}
Expr::Extract { source, .. } => self.resolve_expr_subqueries(source, cancel)?,
Expr::WindowFunction {
args,
partition_by,
order_by,
..
} => {
for a in args {
self.resolve_expr_subqueries(a, cancel)?;
}
for p in partition_by {
self.resolve_expr_subqueries(p, cancel)?;
}
for (e, _, _) in order_by {
self.resolve_expr_subqueries(e, cancel)?;
}
}
Expr::ScalarSubquery(_)
| Expr::Exists { .. }
| Expr::InSubquery { .. }
| Expr::RowInSubquery { .. }
| Expr::RowCmpSubquery { .. }
| Expr::Literal(_)
| Expr::Placeholder(_)
| Expr::Column(_) => {}
Expr::InList { expr, list, .. } => {
self.resolve_expr_subqueries(expr, cancel)?;
for item in list {
self.resolve_expr_subqueries(item, cancel)?;
}
}
Expr::Array(items) => {
for elem in items {
self.resolve_expr_subqueries(elem, cancel)?;
}
}
Expr::ArraySubscript { target, index } => {
self.resolve_expr_subqueries(target, cancel)?;
self.resolve_expr_subqueries(index, cancel)?;
}
Expr::ArraySlice { target, lo, hi } => {
self.resolve_expr_subqueries(target, cancel)?;
if let Some(l) = lo {
self.resolve_expr_subqueries(l, cancel)?;
}
if let Some(h) = hi {
self.resolve_expr_subqueries(h, cancel)?;
}
}
Expr::AnyAll { expr, array, .. } => {
self.resolve_expr_subqueries(expr, cancel)?;
if let Expr::ScalarSubquery(inner) = array.as_mut() {
if !crate::subquery::select_is_correlated(inner) {
let s = (**inner).clone();
**array = self.materialize_quantified_rows(&s, cancel)?;
}
} else {
self.resolve_expr_subqueries(array, cancel)?;
}
}
Expr::Case {
operand,
branches,
else_branch,
} => {
if let Some(o) = operand {
self.resolve_expr_subqueries(o, cancel)?;
}
for (w, t) in branches {
self.resolve_expr_subqueries(w, cancel)?;
self.resolve_expr_subqueries(t, cancel)?;
}
if let Some(e) = else_branch {
self.resolve_expr_subqueries(e, cancel)?;
}
}
}
Ok(())
}
}
impl Engine {
pub(crate) fn project_row_simple(
&self,
row: &Row<'static>,
items: &[SelectItem],
schema_cols: &[ColumnSchema],
alias: &str,
) -> Result<Row<'static>, EngineError> {
let ctx = self.ev_ctx(schema_cols, Some(alias));
let cancel = CancelToken::none();
let mut out_vals = Vec::new();
for item in items {
match item {
SelectItem::Wildcard | SelectItem::QualifiedWildcard(_) => {
out_vals.extend(row.values.iter().cloned());
}
SelectItem::Expr { expr, .. } => {
let v = self.eval_expr_with_correlated(expr, row, &ctx, cancel, None)?;
out_vals.push(v);
}
}
}
Ok(Row::new(out_vals))
}
pub(crate) fn derive_output_columns(
&self,
items: &[SelectItem],
schema_cols: &[ColumnSchema],
table_alias: &str,
) -> Vec<ColumnSchema> {
let mut out = Vec::new();
for item in items {
match item {
SelectItem::Wildcard | SelectItem::QualifiedWildcard(_) => {
out.extend(schema_cols.iter().cloned());
}
SelectItem::Expr { expr, alias } => {
if let Expr::Column(col) = expr
&& let Some(sc) = schema_cols.iter().find(|c| c.name == col.name)
{
let name = alias.clone().unwrap_or_else(|| sc.name.clone());
let mut c = ColumnSchema::new(name, sc.ty, sc.nullable);
c.user_enum_type = sc.user_enum_type.clone();
out.push(c);
continue;
}
let name = alias.clone().unwrap_or_else(|| "?column?".to_string());
let (ty, nullable) = build_projection(
core::slice::from_ref(item),
schema_cols,
table_alias,
self.backslash_escapes,
)
.ok()
.and_then(|p| p.into_iter().next())
.map_or((DataType::Text, true), |p| (p.ty, p.nullable));
out.push(ColumnSchema::new(name, ty, nullable));
}
}
}
out
}
fn meta_view_result(&self, name: &str) -> Option<QueryResult> {
Some(match name {
"spg_statistic" => self.exec_spg_statistic(),
"spg_stat_replication" => self.exec_spg_stat_replication(),
"spg_stat_segment" => self.exec_spg_stat_segment(),
"spg_memory_stats" => self.exec_spg_memory_stats(),
"spg_stat_query" => self.exec_spg_stat_query(),
"pg_stat_statements" => self.exec_pg_stat_statements(),
"spg_stat_activity" => self.exec_spg_stat_activity(),
"pg_stat_activity" => self.exec_pg_stat_activity(),
"pg_locks" => self.exec_pg_locks(),
"pg_statio_user_tables" => self.exec_pg_statio_user_tables(),
"spg_stat_mvcc" => self.exec_spg_stat_mvcc(),
"spg_partition_health" => self.exec_spg_partition_health(),
"spg_audit_chain" => self.exec_spg_audit_chain(),
"spg_audit_verify" => self.exec_spg_audit_verify(),
"spg_table_ddl" => self.exec_spg_table_ddl(),
"spg_role_ddl" => self.exec_spg_role_ddl(),
"spg_database_ddl" => self.exec_spg_database_ddl(),
_ => return None,
})
}
pub(crate) fn admin_view_catalog(&self, stmt: &SelectStatement) -> Option<Catalog> {
let from = stmt.from.as_ref()?;
if !from.joins.is_empty() || self.active_catalog().get(&from.primary.name).is_some() {
return None;
}
let lower = from.primary.name.to_ascii_lowercase();
let QueryResult::Rows { columns, rows } = self.meta_view_result(&lower)? else {
return None;
};
let mut catalog = self.active_catalog().clone();
let cols = infer_column_types(&columns, &rows);
catalog
.create_table(TableSchema::new(from.primary.name.clone(), cols))
.ok()?;
Some(catalog)
}
pub(crate) fn exec_select_cancel(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
self.exec_select_cancel_as(stmt, cancel, None)
}
fn try_bare_count_star(
&self,
stmt: &SelectStatement,
as_role: Option<&str>,
) -> Result<Option<QueryResult>, EngineError> {
use spg_sql::ast::SelectItem;
if as_role.is_some()
|| !stmt.ctes.is_empty()
|| !stmt.unions.is_empty()
|| stmt.where_.is_some()
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| stmt.distinct
|| !stmt.order_by.is_empty()
|| stmt.limit.is_some()
|| stmt.offset.is_some()
|| stmt.items.len() != 1
{
return Ok(None);
}
let Some(from) = &stmt.from else {
return Ok(None);
};
if !from.joins.is_empty()
|| stmt.locking.is_some()
|| from.primary.lateral_subquery.is_some()
|| from.primary.unnest_expr.is_some()
|| from.primary.generate_series_args.is_some()
|| from.primary.name.is_empty()
|| from.primary.name.starts_with("__spg_")
{
return Ok(None);
}
if crate::partition::has_children(self.active_catalog(), &from.primary.name) {
return Ok(None);
}
let SelectItem::Expr { expr, alias } = &stmt.items[0] else {
return Ok(None);
};
let spg_sql::ast::Expr::FunctionCall { name, args } = expr else {
return Ok(None);
};
if !name.eq_ignore_ascii_case("count_star") || !args.is_empty() {
return Ok(None);
}
let Some(table) = self.active_catalog().get(&from.primary.name) else {
return Ok(None);
};
if table.schema().row_security {
return Ok(None);
}
if table.has_cold_rows_fast() {
return Ok(None);
}
let n = table.count_visible(&self.current_snapshot());
let col = alias.clone().unwrap_or_else(|| String::from("count"));
Ok(Some(QueryResult::Rows {
columns: alloc::vec![ColumnSchema::new(col, DataType::BigInt, false)],
rows: alloc::vec![Row::new(alloc::vec![Value::BigInt(
i64::try_from(n).unwrap_or(i64::MAX)
)])],
}))
}
pub(crate) fn index_only_shape<'s>(
&'s self,
stmt: &'s SelectStatement,
) -> Option<(&'s spg_storage::Table, &'s str, usize, String)> {
use spg_sql::ast::SelectItem;
if !stmt.ctes.is_empty()
|| !stmt.unions.is_empty()
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| stmt.distinct
|| stmt.locking.is_some()
|| !stmt.order_by.is_empty()
|| stmt.limit.is_some()
|| stmt.offset.is_some()
|| stmt.items.len() != 1
{
return None;
}
let (Some(from), Some(_)) = (&stmt.from, &stmt.where_) else {
return None;
};
if !from.joins.is_empty()
|| from.primary.lateral_subquery.is_some()
|| from.primary.unnest_expr.is_some()
|| from.primary.generate_series_args.is_some()
|| from.primary.name.is_empty()
|| from.primary.name.starts_with("__spg_")
{
return None;
}
if crate::partition::has_children(self.active_catalog(), &from.primary.name) {
return None;
}
let SelectItem::Expr { expr, alias } = &stmt.items[0] else {
return None;
};
let spg_sql::ast::Expr::Column(c) = expr else {
return None;
};
let alias_name = from.primary.alias.as_deref().unwrap_or(&from.primary.name);
if let Some(q) = c.qualifier.as_deref()
&& !q.eq_ignore_ascii_case(alias_name)
{
return None;
}
let table = self.active_catalog().get(&from.primary.name)?;
if table.schema().row_security {
return None;
}
let cols = &table.schema().columns;
let pos = cols
.iter()
.position(|s| s.name.eq_ignore_ascii_case(&c.name))?;
let out = alias.clone().unwrap_or_else(|| cols[pos].name.clone());
Some((table, alias_name, pos, out))
}
pub(crate) fn stmt_takes_index_only_scan(&self, stmt: &SelectStatement) -> bool {
let Some((table, alias_name, pos, _)) = self.index_only_shape(stmt) else {
return false;
};
let Some(where_) = stmt.where_.as_ref() else {
return false;
};
crate::index_access::index_only_precheck(
where_,
&table.schema().columns,
table,
alias_name,
pos,
)
.is_some()
}
fn try_index_only_scan(
&self,
stmt: &SelectStatement,
) -> Result<Option<QueryResult>, EngineError> {
let Some((table, alias_name, pos, out_name)) = self.index_only_shape(stmt) else {
return Ok(None);
};
let where_ = stmt.where_.as_ref().expect("shape checked it");
let cols = &table.schema().columns;
let Some(values) = crate::index_access::try_index_only_range(
where_,
cols,
table,
alias_name,
&self.current_snapshot(),
pos,
) else {
return Ok(None);
};
let schema = alloc::vec![ColumnSchema::new(
out_name,
cols[pos].ty,
cols[pos].nullable
)];
Ok(Some(QueryResult::Rows {
columns: schema,
rows: values
.into_iter()
.map(|v| Row::new(alloc::vec![v]))
.collect(),
}))
}
pub(crate) fn try_index_only_stream<F>(
&self,
stmt: &SelectStatement,
emit: &mut F,
) -> Result<Option<usize>, EngineError>
where
F: FnMut(crate::StreamItem<'_>) -> Result<(), EngineError>,
{
let Some((table, alias_name, pos, out_name)) = self.index_only_shape(stmt) else {
return Ok(None);
};
let where_ = stmt.where_.as_ref().expect("shape checked it");
let cols = &table.schema().columns;
let schema = alloc::vec![ColumnSchema::new(
out_name,
cols[pos].ty,
cols[pos].nullable
)];
let snapshot = self.current_snapshot();
let mut wrote_header = false;
let counted = crate::index_access::index_only_range_each(
where_,
cols,
table,
alias_name,
&snapshot,
pos,
&mut |v: spg_storage::Value<'_>| {
if !wrote_header {
emit(crate::StreamItem::Header(&schema))?;
wrote_header = true;
}
emit(crate::StreamItem::Row(crate::RowCells::Refs(&[&v])))
},
);
match counted {
None => Ok(None),
Some(Err(e)) => Err(e),
Some(Ok(n)) => {
if !wrote_header {
emit(crate::StreamItem::Header(&schema))?;
}
Ok(Some(n))
}
}
}
#[inline(never)]
fn apply_distinct_on(
&self,
result: QueryResult,
don_hidden: usize,
don_limit: &(
Option<spg_sql::ast::LimitExpr>,
Option<spg_sql::ast::LimitExpr>,
),
don_top1: usize,
orig_order_by: &[spg_sql::ast::OrderBy],
) -> Result<QueryResult, EngineError> {
let QueryResult::Rows { columns, rows } = result else {
return Ok(result);
};
let mut kept: alloc::vec::Vec<Row<'static>>;
let key_start;
if don_top1 > 0 {
let tail = don_top1 - 1;
key_start = columns.len().saturating_sub(don_hidden + tail);
let ord_start = key_start + don_hidden;
let tail_dirs: alloc::vec::Vec<(bool, Option<bool>)> = orig_order_by[don_hidden..]
.iter()
.map(|o| (o.desc, o.nulls_first))
.collect();
let mysql = self.backslash_escapes;
let better = |a: &Row<'static>, b: &Row<'static>| -> bool {
for (k, (desc, nf)) in tail_dirs.iter().enumerate() {
let av = a.values.get(ord_start + k).unwrap_or(&Value::Null);
let bv = b.values.get(ord_start + k).unwrap_or(&Value::Null);
match crate::order_by_value_cmp_in(*desc, *nf, av, bv, mysql) {
core::cmp::Ordering::Less => return true,
core::cmp::Ordering::Greater => return false,
core::cmp::Ordering::Equal => {}
}
}
false
};
let mut slot: hashbrown::HashMap<String, usize> = hashbrown::HashMap::new();
let mut best: alloc::vec::Vec<Row<'static>> = alloc::vec::Vec::new();
let mut keybuf = String::new();
for row in rows {
keybuf.clear();
for v in row.values.get(key_start..ord_start).unwrap_or(&[]) {
aggregate::push_canonical_key(&mut keybuf, v);
}
match slot.get(keybuf.as_str()) {
Some(&i) => {
if better(&row, &best[i]) {
best[i] = row;
}
}
None => {
slot.insert(keybuf.clone(), best.len());
best.push(row);
}
}
}
let full_dirs: alloc::vec::Vec<(bool, Option<bool>)> = orig_order_by
.iter()
.map(|o| (o.desc, o.nulls_first))
.collect();
best.sort_by(|a, b| {
for (k, (desc, nf)) in full_dirs.iter().enumerate() {
let av = a.values.get(key_start + k).unwrap_or(&Value::Null);
let bv = b.values.get(key_start + k).unwrap_or(&Value::Null);
match crate::order_by_value_cmp_in(*desc, *nf, av, bv, mysql) {
core::cmp::Ordering::Equal => {}
o => return o,
}
}
core::cmp::Ordering::Equal
});
for r in &mut best {
r.values.truncate(key_start);
}
kept = best;
} else {
key_start = columns.len().saturating_sub(don_hidden);
let mut seen: alloc::vec::Vec<alloc::vec::Vec<Value<'static>>> = alloc::vec::Vec::new();
kept = alloc::vec::Vec::new();
for mut row in rows {
let key: alloc::vec::Vec<Value<'static>> =
row.values.get(key_start..).unwrap_or(&[]).to_vec();
if seen.iter().any(|k| k == &key) {
continue;
}
seen.push(key);
row.values.truncate(key_start);
kept.push(row);
}
}
let mut columns = columns;
columns.truncate(key_start);
let kept = apply_deferred_limit(kept, don_limit);
Ok(QueryResult::Rows {
columns,
rows: kept,
})
}
pub(crate) fn exec_select_cancel_as(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
as_role: Option<&str>,
) -> Result<QueryResult, EngineError> {
if let Some(expanded) = self.expand_aggregate_wildcard(stmt) {
return self.exec_select_cancel_as(&expanded, cancel, as_role);
}
let aliased;
let stmt = if crate::orderby::order_by_names_an_alias(stmt) {
let mut s = stmt.clone();
crate::orderby::resolve_order_by_position(&mut s);
aliased = s;
&aliased
} else {
stmt
};
let don_stmt;
let orig_order_by = stmt.order_by.clone();
let (stmt, don_hidden, don_limit, don_top1) = if stmt.distinct_on.is_empty() {
(stmt, 0, (None, None), 0usize)
} else {
let mut s = stmt.clone();
let hidden = s.distinct_on.len();
for (i, e) in stmt.distinct_on.iter().enumerate() {
s.items.push(SelectItem::Expr {
expr: e.clone(),
alias: Some(alloc::format!("__distinct_on_{i}")),
});
}
let prefix_matches = s.order_by.len() >= hidden
&& stmt
.distinct_on
.iter()
.zip(s.order_by.iter())
.all(|(d, o)| *d == o.expr && !o.desc && o.nulls_first.is_none());
let colls_plain =
crate::orderby::order_by_collations(&s.order_by, &self.ev_ctx(&[], None))
.map(|cs| cs.iter().all(Option::is_none))
.unwrap_or(false);
let top1_tail = if prefix_matches && colls_plain && s.group_by.is_none() {
let tail = s.order_by.len() - hidden;
for (j, o) in s.order_by[hidden..].iter().enumerate() {
s.items.push(SelectItem::Expr {
expr: o.expr.clone(),
alias: Some(alloc::format!("__don_ord_{j}")),
});
}
s.order_by = Vec::new();
tail + 1 } else {
0
};
let deferrable = matches!(
(&s.limit, &s.offset),
(
None | Some(spg_sql::ast::LimitExpr::Literal(_)),
None | Some(spg_sql::ast::LimitExpr::Literal(_))
)
);
let deferred = if deferrable {
(s.limit.take(), s.offset.take())
} else {
(None, None)
};
don_stmt = s;
(&don_stmt, hidden, deferred, top1_tail)
};
self.acl_check_select_as(stmt, as_role)?;
validate_aggregate_placement(stmt)?;
if let Some(r) = self.try_bare_count_star(stmt, as_role)? {
return Ok(r);
}
if let Some(r) = self.try_index_only_scan(stmt)? {
return Ok(r);
}
validate_locking_clause(stmt)?;
let result = self.exec_select_cancel_inner(stmt, cancel)?;
let result = strip_synthetic_order_cols(result);
if stmt.distinct_on.is_empty() {
return Ok(result);
}
self.apply_distinct_on(result, don_hidden, &don_limit, don_top1, &orig_order_by)
}
#[inline(never)]
fn exec_union_chain(
&self,
stmt_ref: &SelectStatement,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
crate::orderby::check_order_by_positions(stmt_ref)?;
let mut head_unknown = branch_unknown_mask(stmt_ref);
let mut head = stmt_ref.clone();
head.unions = Vec::new();
head.order_by = Vec::new();
head.limit = None;
let QueryResult::Rows {
mut columns,
mut rows,
} = self.exec_bare_select_cancel(&head, cancel)?
else {
unreachable!("bare SELECT cannot return CommandOk")
};
for (kind, peer) in &stmt_ref.unions {
let peer_result = if peer.unions.is_empty() {
self.exec_bare_select_cancel(peer, cancel)?
} else {
self.exec_select_cancel(peer, cancel)?
};
let QueryResult::Rows {
columns: peer_cols,
rows: mut peer_rows,
} = peer_result
else {
unreachable!("bare SELECT cannot return CommandOk")
};
if peer_cols.len() != columns.len() {
return Err(EngineError::Unsupported(alloc::format!(
"each {} query must have the same number of columns",
set_op_name(*kind)
)));
}
let peer_unknown = branch_unknown_mask(peer);
for i in 0..columns.len() {
let hu = head_unknown.get(i).copied().unwrap_or(false);
let pu = peer_unknown.get(i).copied().unwrap_or(false);
let (ht, pt) = (columns[i].ty, peer_cols[i].ty);
match (hu, pu) {
(false, false) => {
if !crate::conversions::types_unify(ht, pt) {
return Err(EngineError::Unsupported(alloc::format!(
"{} types {} and {} cannot be matched",
set_op_name(*kind),
crate::conversions::pg_type_name_for_error(ht),
crate::conversions::pg_type_name_for_error(pt),
)));
}
}
(true, false) => {
coerce_branch_column(&mut rows, i, pt, &columns[i].name)?;
columns[i].ty = pt;
head_unknown[i] = false;
}
(false, true) => {
coerce_branch_column(&mut peer_rows, i, ht, &columns[i].name)?;
}
(true, true) => {}
}
}
for (i, pc) in peer_cols.iter().enumerate() {
if pc.nullable {
columns[i].nullable = true;
}
}
let mysql = self.backslash_escapes;
match kind {
UnionKind::All => rows.extend(peer_rows),
UnionKind::Distinct => {
rows.extend(peer_rows);
rows = dedup_rows(rows, mysql);
}
UnionKind::Intersect => {
let idx = PeerIndex::build(&peer_rows, mysql);
rows = dedup_rows(rows, mysql)
.into_iter()
.filter(|r| idx.contains(r))
.collect();
}
UnionKind::IntersectAll => {
let mut idx = PeerIndex::build(&peer_rows, mysql);
let mut kept: Vec<Row<'static>> = Vec::new();
for r in rows {
if idx.take_one(&r) {
kept.push(r);
}
}
rows = kept;
}
UnionKind::Except => {
let idx = PeerIndex::build(&peer_rows, mysql);
rows = dedup_rows(rows, mysql)
.into_iter()
.filter(|r| !idx.contains(r))
.collect();
}
UnionKind::ExceptAll => {
let mut idx = PeerIndex::build(&peer_rows, mysql);
let mut kept: Vec<Row<'static>> = Vec::new();
for r in rows {
if !idx.take_one(&r) {
kept.push(r);
}
}
rows = kept;
}
}
}
unify_union_columns(&mut columns, &mut rows);
if !stmt.order_by.is_empty() {
let synth_ctx = EvalContext::new(&columns, None).with_catalog(self.active_catalog());
let resolved_order: Vec<spg_sql::ast::OrderBy> = stmt
.order_by
.iter()
.map(|o| {
let mut o = o.clone();
if let Expr::Literal(spg_sql::ast::Literal::Integer(n)) = &o.expr
&& *n >= 1
&& let Ok(idx) = usize::try_from(*n - 1)
&& idx < columns.len()
{
o.expr = Expr::Column(spg_sql::ast::ColumnName {
qualifier: None,
name: columns[idx].name.clone(),
});
}
o
})
.collect();
let descs: Vec<bool> = resolved_order.iter().map(|o| o.desc).collect();
let mut tagged: Vec<(Vec<OrderKey>, Row)> = Vec::with_capacity(rows.len());
for r in rows {
let keys = build_order_keys(&resolved_order, &r, &synth_ctx)?;
tagged.push((keys, r));
}
sort_by_keys(&mut tagged, &descs);
rows = tagged.into_iter().map(|(_, r)| r).collect();
}
apply_offset_and_limit(&mut rows, stmt.offset_literal(), stmt.limit_literal());
Ok(QueryResult::Rows { columns, rows })
}
fn exec_select_cancel_inner(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
cancel.check()?;
crate::injection_point!("planner_first_row_fetch", &stmt.from);
if !stmt.window_check_exprs.is_empty() {
let mut probe = stmt.clone();
probe.items = stmt
.window_check_exprs
.iter()
.map(|e| spg_sql::ast::SelectItem::Expr {
expr: e.clone(),
alias: None,
})
.collect();
probe.window_check_exprs = Vec::new();
probe.distinct = false;
probe.distinct_on = Vec::new();
probe.group_by = None;
probe.group_by_all = false;
probe.having = None;
probe.unions = Vec::new();
probe.order_by = Vec::new();
probe.locking = None;
probe.limit = Some(spg_sql::ast::LimitExpr::Literal(0));
probe.offset = None;
probe.limit_with_ties = false;
self.exec_select_cancel_inner(&probe, cancel)?;
}
if let Some(lowered) = self.lower_record_expansion(stmt)? {
return self.exec_select_cancel_inner(&lowered, cancel);
}
if !self.active_catalog().views_all().is_empty() {
if let Some(rewritten) = self.expand_views_in_select(stmt)? {
return self.exec_select_cancel(&rewritten, cancel);
}
}
if let Some(rewritten) = self.expand_partition_parents_in_select(stmt)? {
return self.exec_select_cancel(&rewritten, cancel);
}
if !self.meta_views_materialised && select_references_meta_view(stmt) {
return self.exec_select_with_meta_views(stmt, cancel);
}
if let Some(from) = &stmt.from
&& let Some(seg_id) = from.primary.as_of_segment
{
return self.exec_select_as_of_segment(stmt, from, seg_id);
}
if let Some(from) = &stmt.from
&& from.joins.is_empty()
&& self.active_catalog().get(&from.primary.name).is_none()
{
let lower = from.primary.name.to_ascii_lowercase();
if let Some(result) = self.meta_view_result(&lower) {
let bare = stmt.where_.is_none()
&& stmt.group_by.is_none()
&& stmt.having.is_none()
&& stmt.unions.is_empty()
&& stmt.order_by.is_empty()
&& stmt.limit.is_none()
&& stmt.offset.is_none()
&& !stmt.distinct
&& stmt.items.iter().all(|i| matches!(i, SelectItem::Wildcard));
if bare {
return Ok(result);
}
if let QueryResult::Rows { columns, rows } = result {
let mut catalog = self.active_catalog().clone();
let cols = infer_column_types(&columns, &rows);
let schema = TableSchema::new(from.primary.name.clone(), cols);
catalog.create_table(schema).map_err(EngineError::Storage)?;
let t = catalog
.get_mut(&from.primary.name)
.expect("just-created meta-view table must exist");
for row in rows {
t.insert(row).map_err(EngineError::Storage)?;
}
let mut eng = Engine::restore(catalog);
if let Some(c) = self.clock {
eng = eng.with_clock(c);
}
if let Some(f) = self.salt_fn {
eng = eng.with_salt_fn(f);
}
if let Some(f) = self.backend_pid_fn {
eng.set_backend_pid_fn(f);
}
return eng.exec_select_cancel(stmt, cancel);
}
return Ok(result);
}
}
if !stmt.ctes.is_empty() {
return self.exec_with_ctes(stmt, cancel);
}
let mut stmt_owned;
let stmt_ref: &SelectStatement = if expr_tree_has_subquery(stmt) {
stmt_owned = stmt.clone();
self.pull_up_unique_correlated_agg_subqueries(&mut stmt_owned);
self.pull_up_correlated_limit_one_subqueries(&mut stmt_owned);
self.pull_up_exists_sublinks(&mut stmt_owned);
if !stmt_owned.ctes.is_empty() {
return self.exec_with_ctes(&stmt_owned, cancel);
}
if let Some(out) = self.try_count_star_pk_in_subquery_fast(&stmt_owned, cancel)? {
return Ok(out);
}
self.resolve_select_subqueries(&mut stmt_owned, cancel)?;
&stmt_owned
} else {
stmt
};
if stmt_ref.unions.is_empty() {
return self.exec_bare_select_cancel(stmt_ref, cancel);
}
self.exec_union_chain(stmt_ref, stmt, cancel)
}
#[allow(clippy::too_many_lines)]
#[allow(clippy::too_many_lines)] fn exec_select_unnest(
&self,
stmt: &SelectStatement,
primary: &TableRef,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
let expr = primary
.unnest_expr
.as_deref()
.expect("caller guards unnest_expr.is_some()");
let multi: Option<(alloc::vec::Vec<DataType>, alloc::vec::Vec<Row<'static>>)> =
match unnest_zip_args(expr) {
Some(args) => Some(unnest_zip_rows(args)?),
None => None,
};
let empty_schema: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let ctx = EvalContext::new(&empty_schema, None).with_catalog(self.active_catalog());
let dummy_row = Row::new(alloc::vec::Vec::new());
let mut composite_names: Option<&[&str]> = None;
let (dtypes, rows): (alloc::vec::Vec<DataType>, alloc::vec::Vec<Row<'static>>) =
if let Some(m) = multi {
m
} else {
let unnest_src = {
let v = eval::eval_expr(expr, &dummy_row, &ctx).map_err(EngineError::Eval)?;
crate::eval::values::flatten_2d(&v).unwrap_or(v)
};
let mut return_multi: Option<(
alloc::vec::Vec<DataType>,
alloc::vec::Vec<Row<'static>>,
)> = None;
let (elem_dtype, rows): (DataType, alloc::vec::Vec<Row<'static>>) = match unnest_src
{
Value::Null => (DataType::Text, alloc::vec::Vec::new()),
Value::TextArray(items) => {
let rows = items
.into_iter()
.map(|item| {
Row::new(alloc::vec![match item {
Some(s) => Value::text(s),
None => Value::Null,
}])
})
.collect();
(DataType::Text, rows)
}
Value::IntArray(items) => {
let rows = items
.into_iter()
.map(|item| {
Row::new(alloc::vec![match item {
Some(n) => Value::Int(n),
None => Value::Null,
}])
})
.collect();
(DataType::Int, rows)
}
Value::BigIntArray(items) => {
let rows = items
.into_iter()
.map(|item| {
Row::new(alloc::vec![match item {
Some(n) => Value::BigInt(n),
None => Value::Null,
}])
})
.collect();
(DataType::BigInt, rows)
}
Value::Multirange { kind, ranges } => {
let rows = ranges
.iter()
.map(|sp| {
Row::new(alloc::vec![Value::Range {
kind,
lower: sp.lower.clone(),
upper: sp.upper.clone(),
lower_inc: sp.lower_inc,
upper_inc: sp.upper_inc,
empty: false,
}])
})
.collect();
(DataType::Range(kind), rows)
}
Value::TsVector(lexemes) => {
composite_names = Some(&["lexeme", "positions", "weights"]);
let rows = lexemes
.iter()
.map(|l| {
let (pos, wts) = if l.positions.is_empty() {
(Value::Null, Value::Null)
} else {
let letter = match l.weight {
3 => "A",
2 => "B",
1 => "C",
_ => "D",
};
(
Value::SmallIntArray(
l.positions
.iter()
.map(|p| {
Some(i16::try_from(*p).unwrap_or(i16::MAX))
})
.collect(),
),
Value::TextArray(
l.positions
.iter()
.map(|_| Some(letter.into()))
.collect(),
),
)
};
Row::new(alloc::vec![Value::text(l.word.clone()), pos, wts])
})
.collect();
return_multi = Some((
alloc::vec![
DataType::Text,
DataType::SmallIntArray,
DataType::TextArray
],
rows,
));
(DataType::Text, alloc::vec::Vec::new())
}
other => {
return Err(EngineError::Eval(EvalError::TypeMismatch {
detail: alloc::format!(
"unnest() expects an array argument, got {}",
crate::conversions::pg_type_name_for_error_opt(other.data_type())
),
}));
}
};
if let Some(m) = return_multi {
m
} else {
(alloc::vec![elem_dtype], rows)
}
};
let alias = primary
.alias
.clone()
.unwrap_or_else(|| "unnest".to_string());
let n_vals = dtypes.len();
let mut schema_cols: alloc::vec::Vec<ColumnSchema> = dtypes
.iter()
.enumerate()
.map(|(i, dt)| {
let name = primary
.unnest_column_aliases
.get(i)
.cloned()
.unwrap_or_else(|| {
if let Some(names) = composite_names {
names
.get(i)
.map_or_else(|| "unnest".to_string(), |n| (*n).to_string())
} else if n_vals == 1 {
alias.clone()
} else {
"unnest".to_string()
}
});
ColumnSchema::new(name, *dt, true)
})
.collect();
if primary.scalar_fn_item && schema_cols.len() == 1 {
schema_cols[0].scalar_row_source = true;
}
let rows = if primary.with_ordinality {
let ord_name = primary
.unnest_column_aliases
.get(n_vals)
.cloned()
.unwrap_or_else(|| "ordinality".to_string());
schema_cols.push(ColumnSchema::new(ord_name, DataType::BigInt, false));
rows.into_iter()
.enumerate()
.map(|(i, row)| {
let mut vals = row.values.clone();
vals.push(Value::BigInt(i as i64 + 1));
Row::new(vals)
})
.collect()
} else {
rows
};
let scan_ctx = self.ev_ctx(&schema_cols, Some(&alias));
let filtered: alloc::vec::Vec<Row<'static>> = if let Some(w) = &stmt.where_ {
let mut out = alloc::vec::Vec::with_capacity(rows.len());
for row in rows {
cancel.check()?;
let v = eval::eval_expr(w, &row, &scan_ctx).map_err(EngineError::Eval)?;
if matches!(v, Value::Bool(true)) {
out.push(row);
}
}
out
} else {
rows
};
if aggregate::uses_aggregate(stmt) {
let agg_memo = core::cell::RefCell::new(memoize::MemoizeCache::default());
let agg_correlated = |e: &Expr, r: &Row<'static>, c: &EvalContext<'_>| {
self.eval_expr_with_correlated(e, r, c, cancel, Some(&mut agg_memo.borrow_mut()))
.map_err(|err| match err {
EngineError::Eval(ev) => ev,
other => eval::EvalError::TypeMismatch {
detail: alloc::format!("{other}"),
},
})
};
let agg = aggregate::run(
stmt,
crate::join::AggRows::Owned(&filtered),
&schema_cols,
Some(&alias),
Some(&agg_correlated),
self.parallel_runner.0.as_deref(),
Some(self.active_catalog()),
Some(self),
)?;
return self.finish_agg_result(agg, stmt, cancel);
}
let projection =
build_projection(&stmt.items, &schema_cols, &alias, self.backslash_escapes)?;
let mut projected_rows: alloc::vec::Vec<Row<'static>> =
alloc::vec::Vec::with_capacity(filtered.len());
let srf_idxs = self.srf_target_idxs(&projection);
let mut src_of_row: alloc::vec::Vec<usize> = alloc::vec::Vec::new();
if !srf_idxs.is_empty() {
let (rows, src) =
expand_projection_srfs(self, &projection, &srf_idxs, &filtered, &scan_ctx)?;
projected_rows = rows;
src_of_row = src;
} else {
let mut proj_memo = memoize::MemoizeCache::default();
for row in &filtered {
let mut vals = alloc::vec::Vec::with_capacity(projection.len());
for p in &projection {
vals.push(self.eval_expr_with_correlated(
&p.expr,
row,
&scan_ctx,
cancel,
Some(&mut proj_memo),
)?);
}
projected_rows.push(Row::new(vals));
}
}
let columns: alloc::vec::Vec<ColumnSchema> = projection
.iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name.clone(), p.ty, p.nullable);
c.user_enum_type = p.user_enum_type.clone();
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
let order_by = resolve_positional_order_by(&stmt.order_by, &projection);
if !order_by.is_empty() {
let out_cols = if srf_idxs.is_empty() {
alloc::vec![None; order_by.len()]
} else {
srf_order_output_cols(&order_by, &projection)
};
let mut indexed: alloc::vec::Vec<(usize, Vec<Value<'static>>)> = projected_rows
.iter()
.enumerate()
.map(|(k, out)| -> Result<_, EngineError> {
let src = src_of_row.get(k).copied().unwrap_or(k);
let keys: Result<Vec<Value<'static>>, EngineError> = order_by
.iter()
.zip(out_cols.iter())
.map(|(ob, oc)| srf_order_key(ob, *oc, out, &filtered[src], &scan_ctx))
.collect();
Ok((k, keys?))
})
.collect::<Result<_, _>>()?;
indexed.sort_by(|a, b| {
for (idx, (ka, kb)) in a.1.iter().zip(b.1.iter()).enumerate() {
let o = &order_by[idx];
let cmp = order_by_value_cmp_in(
o.desc,
o.nulls_first,
ka,
kb,
scan_ctx.mysql_dialect && !crate::eval::is_binary_coerced(&o.expr),
);
if cmp != core::cmp::Ordering::Equal {
return cmp;
}
}
core::cmp::Ordering::Equal
});
projected_rows = indexed
.into_iter()
.map(|(i, _)| projected_rows[i].clone())
.collect();
}
if stmt.distinct {
projected_rows = dedup_rows(projected_rows, scan_ctx.mysql_dialect);
}
if let Some(offset) = stmt.offset_literal() {
let off = (offset as usize).min(projected_rows.len());
projected_rows.drain(..off);
}
if let Some(limit) = stmt.limit_literal() {
projected_rows.truncate(limit as usize);
}
Ok(QueryResult::Rows {
columns,
rows: projected_rows,
})
}
fn exec_select_generate_series(
&self,
stmt: &SelectStatement,
primary: &TableRef,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
let args = primary
.generate_series_args
.as_ref()
.expect("caller guards generate_series_args.is_some()");
let (elem_dtype, rows) = generate_series_rows(args, &cancel)?;
let alias = primary
.alias
.clone()
.unwrap_or_else(|| "generate_series".to_string());
let col_name = primary
.unnest_column_aliases
.first()
.cloned()
.unwrap_or_else(|| alias.clone());
let col_schema = ColumnSchema::new(col_name, elem_dtype, true);
let mut schema_cols = alloc::vec![col_schema.clone()];
let rows = if primary.with_ordinality {
let ord_name = primary
.unnest_column_aliases
.get(1)
.cloned()
.unwrap_or_else(|| "ordinality".to_string());
schema_cols.push(ColumnSchema::new(ord_name, DataType::BigInt, false));
rows.into_iter()
.enumerate()
.map(|(i, row)| {
let mut vals = row.values.clone();
vals.push(Value::BigInt(i as i64 + 1));
Row::new(vals)
})
.collect()
} else {
rows
};
let scan_ctx = self.ev_ctx(&schema_cols, Some(&alias));
let filtered: alloc::vec::Vec<Row<'static>> = if let Some(w) = &stmt.where_ {
let mut out = alloc::vec::Vec::with_capacity(rows.len());
for row in rows {
cancel.check()?;
let v = eval::eval_expr(w, &row, &scan_ctx).map_err(EngineError::Eval)?;
if matches!(v, Value::Bool(true)) {
out.push(row);
}
}
out
} else {
rows
};
if aggregate::uses_aggregate(stmt) {
let agg_memo = core::cell::RefCell::new(memoize::MemoizeCache::default());
let agg_correlated = |e: &Expr, r: &Row<'static>, c: &EvalContext<'_>| {
self.eval_expr_with_correlated(e, r, c, cancel, Some(&mut agg_memo.borrow_mut()))
.map_err(|err| match err {
EngineError::Eval(ev) => ev,
other => eval::EvalError::TypeMismatch {
detail: alloc::format!("{other}"),
},
})
};
let agg = aggregate::run(
stmt,
crate::join::AggRows::Owned(&filtered),
&schema_cols,
Some(&alias),
Some(&agg_correlated),
self.parallel_runner.0.as_deref(),
Some(self.active_catalog()),
Some(self),
)?;
return self.finish_agg_result(agg, stmt, cancel);
}
let projection =
build_projection(&stmt.items, &schema_cols, &alias, self.backslash_escapes)?;
let srf_idxs = self.srf_target_idxs(&projection);
let mut src_of_row: alloc::vec::Vec<usize> = alloc::vec::Vec::new();
let mut projected_rows: alloc::vec::Vec<Row<'static>> =
alloc::vec::Vec::with_capacity(filtered.len());
let mut proj_memo = memoize::MemoizeCache::default();
if !srf_idxs.is_empty() {
let (rows, src) =
expand_projection_srfs(self, &projection, &srf_idxs, &filtered, &scan_ctx)?;
projected_rows = rows;
src_of_row = src;
} else {
for row in &filtered {
let mut vals = alloc::vec::Vec::with_capacity(projection.len());
for p in &projection {
vals.push(self.eval_expr_with_correlated(
&p.expr,
row,
&scan_ctx,
cancel,
Some(&mut proj_memo),
)?);
}
projected_rows.push(Row::new(vals));
}
}
let columns: alloc::vec::Vec<ColumnSchema> = projection
.iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name.clone(), p.ty, p.nullable);
c.user_enum_type = p.user_enum_type.clone();
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
let order_by = resolve_positional_order_by(&stmt.order_by, &projection);
if !order_by.is_empty() {
let out_cols = if srf_idxs.is_empty() {
alloc::vec![None; order_by.len()]
} else {
srf_order_output_cols(&order_by, &projection)
};
let mut indexed: alloc::vec::Vec<(usize, Vec<Value<'static>>)> = projected_rows
.iter()
.enumerate()
.map(|(k, out)| -> Result<_, EngineError> {
let r = &filtered[src_of_row.get(k).copied().unwrap_or(k)];
let keys: Result<Vec<Value<'static>>, EngineError> = order_by
.iter()
.zip(out_cols.iter())
.map(|(ob, oc)| srf_order_key(ob, *oc, out, r, &scan_ctx))
.collect();
Ok((k, keys?))
})
.collect::<Result<_, _>>()?;
indexed.sort_by(|a, b| {
for (idx, (ka, kb)) in a.1.iter().zip(b.1.iter()).enumerate() {
let o = &stmt.order_by[idx];
let cmp = order_by_value_cmp_in(
o.desc,
o.nulls_first,
ka,
kb,
scan_ctx.mysql_dialect && !crate::eval::is_binary_coerced(&o.expr),
);
if cmp != core::cmp::Ordering::Equal {
return cmp;
}
}
core::cmp::Ordering::Equal
});
projected_rows = indexed
.into_iter()
.map(|(i, _)| projected_rows[i].clone())
.collect();
}
if stmt.distinct {
projected_rows = dedup_rows(projected_rows, scan_ctx.mysql_dialect);
}
if let Some(offset) = stmt.offset_literal() {
let off = (offset as usize).min(projected_rows.len());
projected_rows.drain(..off);
}
if let Some(limit) = stmt.limit_literal() {
projected_rows.truncate(limit as usize);
}
Ok(QueryResult::Rows {
columns,
rows: projected_rows,
})
}
#[inline(never)]
fn try_from_shape_paths(
&self,
stmt: &SelectStatement,
from: &spg_sql::ast::FromClause,
cancel: CancelToken<'_>,
) -> Result<Option<QueryResult>, EngineError> {
if !from.joins.is_empty() {
if let Some(eliminated) = self.try_eliminate_redundant_left_joins(stmt) {
return self.exec_bare_select_cancel(&eliminated, cancel).map(Some);
}
if !self.env_cfg().disable_joinfold {
if let Some(folded) = self.try_fold_inner_joins(stmt, cancel)? {
return self.exec_bare_select_cancel(&folded, cancel).map(Some);
}
}
return self.exec_joined_select(stmt, from, cancel).map(Some);
}
if from.primary.unnest_expr.is_some() {
return self
.exec_select_unnest(stmt, &from.primary, cancel)
.map(Some);
}
if from.primary.jsonb_each_text_arg.is_some() {
return self
.exec_select_jsonb_each_text(stmt, &from.primary, cancel)
.map(Some);
}
if from.primary.rows_from.is_some() {
let (rows, mut schema_cols) = self.rows_from_rows(&from.primary)?;
for (i, new_name) in from.primary.unnest_column_aliases.iter().enumerate() {
if let Some(col) = schema_cols.get_mut(i) {
col.name = new_name.clone();
}
}
let alias = from
.primary
.alias
.clone()
.unwrap_or_else(|| from.primary.name.clone());
return self
.exec_select_over_rows(stmt, rows, schema_cols, &alias, cancel)
.map(Some);
}
if let Some(jt) = &from.primary.json_table {
let (rows, schema_cols) = self.json_table_rows(jt, None)?;
let alias = from
.primary
.alias
.clone()
.unwrap_or_else(|| from.primary.name.clone());
return self
.exec_select_over_rows(stmt, rows, schema_cols, &alias, cancel)
.map(Some);
}
if from.primary.table_fn_call.is_some() {
let (rows, mut schema_cols) = self.table_fn_rows(&from.primary)?;
let rows = if from.primary.with_ordinality {
schema_cols.push(ColumnSchema::new(
"ordinality".to_string(),
DataType::BigInt,
false,
));
rows.into_iter()
.enumerate()
.map(|(i, r)| {
let mut vals = r.values;
vals.push(Value::BigInt(i as i64 + 1));
Row::new(vals)
})
.collect()
} else {
rows
};
for (i, new_name) in from.primary.unnest_column_aliases.iter().enumerate() {
if let Some(col) = schema_cols.get_mut(i) {
col.name = new_name.clone();
}
}
let alias = from
.primary
.alias
.clone()
.unwrap_or_else(|| from.primary.name.clone());
return self
.exec_select_over_rows(stmt, rows, schema_cols, &alias, cancel)
.map(Some);
}
if from.joins.is_empty() && from.primary.lateral_subquery.is_some() {
if let Some(flat) = try_flatten_derived(stmt, &from.primary) {
return self.exec_select_cancel(&flat, cancel).map(Some);
}
if let Some(rewritten) = try_count_over_offset(stmt, &from.primary) {
return self.exec_select_cancel(&rewritten, cancel).map(Some);
}
if let Some(rewritten) = try_count_over_const_unnest(stmt, &from.primary) {
return self.exec_select_cancel(&rewritten, cancel).map(Some);
}
return self
.exec_select_derived(stmt, &from.primary, cancel)
.map(Some);
}
if from.primary.generate_series_args.is_some() {
return self
.exec_select_generate_series(stmt, &from.primary, cancel)
.map(Some);
}
Ok(None)
}
#[inline(never)]
fn pick_indexed_rows<'r>(
&'r self,
stmt: &SelectStatement,
table: &'r spg_storage::Table,
schema_cols: &[spg_storage::ColumnSchema],
alias: &str,
ctx: &crate::eval::EvalContext<'_>,
seek_snapshot: &crate::Snapshot,
) -> Option<Vec<Cow<'r, Row<'static>>>> {
stmt.where_.as_ref().and_then(|w| {
try_index_seek(
w,
schema_cols,
self.active_catalog(),
table,
alias,
seek_snapshot,
)
.or_else(|| {
try_gin_seek(
w,
schema_cols,
self.active_catalog(),
table,
alias,
ctx,
seek_snapshot,
)
})
.or_else(|| {
try_trgm_seek(w, schema_cols, table, alias, seek_snapshot)
})
.or_else(|| {
try_gin_jsonb_seek(w, schema_cols, table, alias, seek_snapshot)
})
})
}
#[inline(never)]
fn try_seek_fast_paths(
&self,
stmt: &SelectStatement,
table: &spg_storage::Table,
schema_cols: &[spg_storage::ColumnSchema],
alias: &str,
seek_snapshot: &crate::Snapshot,
cancel: CancelToken<'_>,
) -> Result<Option<QueryResult>, EngineError> {
if let Some(nsw_rows) = try_nsw_knn(stmt, table, schema_cols, alias, seek_snapshot) {
let ordered: Vec<Cow<'_, Row<'static>>> = nsw_rows
.into_iter()
.filter_map(|i| table.rows().get(i).map(Cow::Borrowed))
.collect();
return materialise_in_order(
stmt,
schema_cols,
alias,
&ordered,
self.backslash_escapes,
)
.map(Some);
}
if let Some(walked) = try_pk_walk_top_n(
stmt,
self.active_catalog(),
table,
schema_cols,
alias,
self,
cancel,
) {
return materialise_in_order(stmt, schema_cols, alias, &walked, self.backslash_escapes)
.map(Some);
}
if aggregate::uses_aggregate(stmt)
&& let Some(out) = self.try_count_star_pk_in_list_fast(stmt, table, schema_cols, alias)
{
return Ok(Some(out));
}
if aggregate::uses_aggregate(stmt)
&& let Some(out) = self.try_count_star_indexed_range_fast(
stmt,
table,
schema_cols,
alias,
seek_snapshot,
)
{
return Ok(Some(out));
}
Ok(None)
}
#[inline(never)]
fn try_pre_from_paths(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<Option<QueryResult>, EngineError> {
if !self.meta_views_materialised && select_references_meta_view(stmt) {
return self.exec_select_with_meta_views(stmt, cancel).map(Some);
}
if select_has_window(stmt) {
if let Some(rewritten) = rewrite_agg_before_window(stmt) {
return self.exec_select_cancel(&rewritten, cancel).map(Some);
}
return self.exec_select_with_window(stmt, cancel).map(Some);
}
Ok(None)
}
#[inline(never)]
fn try_ctid_projection(
&self,
stmt: &SelectStatement,
primary: &spg_sql::ast::TableRef,
table: &spg_storage::Table,
schema_cols: &[spg_storage::ColumnSchema],
alias: &str,
cancel: CancelToken<'_>,
) -> Result<Option<QueryResult>, EngineError> {
if references_ctid(stmt) {
let snapshot = self.current_snapshot();
let mut ext_cols = schema_cols.to_vec();
for name in SYSTEM_COLUMNS {
ext_cols.push(ColumnSchema::new(name.to_string(), DataType::Text, false));
}
let table_oid =
crate::system_catalog::relation_oid(self.active_catalog(), &primary.name)
.unwrap_or(0);
let headers = table.headers();
let rows: Vec<Row<'static>> = table
.scan_visible(&snapshot)
.map(|(i, r)| {
let mut vals = r.values.clone();
vals.push(Value::Tid(0, i as u32 + 1));
let h = headers.get(i);
vals.push(Value::Xid(h.map_or(0, |h| h.xmin as u32)));
vals.push(Value::Xid(h.map_or(0, |h| h.xmax as u32)));
vals.push(Value::Cid(0));
vals.push(Value::Cid(0));
vals.push(Value::BigInt(table_oid));
Row::new(vals)
})
.collect();
return self
.exec_select_over_rows(stmt, rows, ext_cols, alias, cancel)
.map(Some);
}
Ok(None)
}
#[inline(never)]
fn try_sequence_relation(
&self,
stmt: &SelectStatement,
primary: &spg_sql::ast::TableRef,
cancel: CancelToken<'_>,
) -> Result<Option<QueryResult>, EngineError> {
if self.active_catalog().get(&primary.name).is_none()
&& let Some(seq) = self.active_catalog().sequence(&primary.name)
{
let rows = alloc::vec![Row::new(alloc::vec![
Value::BigInt(seq.last_value),
Value::BigInt(0),
Value::Bool(seq.is_called),
])];
let schema_cols = alloc::vec![
ColumnSchema::new("last_value", DataType::BigInt, false),
ColumnSchema::new("log_cnt", DataType::BigInt, false),
ColumnSchema::new("is_called", DataType::Bool, false),
];
let alias = primary
.alias
.clone()
.unwrap_or_else(|| primary.name.clone());
return self
.exec_select_over_rows(stmt, rows, schema_cols, &alias, cancel)
.map(Some);
}
Ok(None)
}
pub(crate) fn exec_bare_select_cancel(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
check_with_ties_requires_order_by(stmt)?;
crate::window::reject_window_in_row_clauses(stmt)?;
crate::orderby::check_order_by_legality(stmt)?;
if let Some(rewritten) = self.desugar_using_natural(stmt)? {
return self.exec_bare_select_cancel(&rewritten, cancel);
}
if let Some(rewritten) = self.rls_rewrite_joins(stmt) {
return self.exec_bare_select_cancel(&rewritten, cancel);
}
let rls_stmt;
let stmt = match self.rls_select_predicate(stmt)? {
Some(pred) => {
let mut s = stmt.clone();
s.where_ = Some(match s.where_.take() {
Some(existing) => spg_sql::ast::Expr::Binary {
lhs: alloc::boxed::Box::new(existing),
op: spg_sql::ast::BinOp::And,
rhs: alloc::boxed::Box::new(pred),
},
None => pred,
});
rls_stmt = s;
&rls_stmt
}
None => stmt,
};
if let Some(done) = self.try_pre_from_paths(stmt, cancel)? {
return Ok(done);
}
let Some(from) = &stmt.from else {
return self.exec_constant_select(stmt);
};
if let Some(done) = self.try_from_shape_paths(stmt, from, cancel)? {
return Ok(done);
}
let primary = &from.primary;
if let Some(done) = self.try_sequence_relation(stmt, primary, cancel)? {
return Ok(done);
}
let table = self.active_catalog().get(&primary.name).ok_or_else(|| {
StorageError::TableNotFound {
name: primary.name.clone(),
}
})?;
let schema_cols = &table.schema().columns;
let alias = primary.alias.as_deref().unwrap_or(primary.name.as_str());
if let Some(done) =
self.try_ctid_projection(stmt, primary, table, schema_cols, alias, cancel)?
{
return Ok(done);
}
let ctx = self.ev_ctx(schema_cols, Some(alias));
let seek_snapshot = self.current_snapshot();
if let Some(done) =
self.try_seek_fast_paths(stmt, table, schema_cols, alias, &seek_snapshot, cancel)?
{
return Ok(done);
}
let indexed_rows =
self.pick_indexed_rows(stmt, table, schema_cols, alias, &ctx, &seek_snapshot);
if aggregate::uses_aggregate(stmt) {
return self.run_single_table_aggregate(
stmt,
table,
schema_cols,
alias,
indexed_rows,
cancel,
);
}
self.run_single_table_scan(stmt, table, schema_cols, alias, indexed_rows, cancel)
}
#[allow(clippy::type_complexity)]
pub(crate) fn json_table_rows(
&self,
jt: &spg_sql::ast::JsonTable,
outer_doc: Option<&crate::json::JsonValue>,
) -> Result<(alloc::vec::Vec<Row<'static>>, alloc::vec::Vec<ColumnSchema>), EngineError> {
let schema = json_table_schema(&jt.columns);
let empty_schema: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let ctx = EvalContext::new(&empty_schema, None);
let dummy = Row::new(alloc::vec::Vec::new());
let vars: Option<crate::json::JsonValue> = if jt.passing.is_empty() {
None
} else {
let mut entries = alloc::vec::Vec::new();
for (name, e) in &jt.passing {
let v = eval::eval_expr(e, &dummy, &ctx).map_err(EngineError::Eval)?;
entries.push((name.clone(), value_to_json_value(&v)));
}
Some(crate::json::JsonValue::Object(entries))
};
let root_owned;
let root: &crate::json::JsonValue = match outer_doc {
Some(d) => d,
None => {
let doc_val = eval::eval_expr(&jt.doc, &dummy, &ctx).map_err(EngineError::Eval)?;
let src = match &doc_val {
Value::Null => return Ok((alloc::vec::Vec::new(), schema)),
Value::Json(s) | Value::Text(s) => s.as_ref().to_string(),
other => {
return Err(EngineError::Unsupported(alloc::format!(
"JSON_TABLE document must be json/text, got {}",
crate::conversions::pg_type_name_for_error_opt(other.data_type())
)));
}
};
root_owned = crate::json::parse_doc(&src).map_err(EngineError::Eval)?;
&root_owned
}
};
let items = crate::json::json_table_path(root, &jt.row_path, vars.as_ref())
.map_err(EngineError::Eval)?;
let mut rows: alloc::vec::Vec<Row<'static>> = alloc::vec::Vec::new();
for (idx, item) in items.iter().enumerate() {
self.json_table_emit_item(jt, item, idx, vars.as_ref(), &mut rows)?;
}
Ok((rows, schema))
}
fn json_table_emit_item(
&self,
jt: &spg_sql::ast::JsonTable,
item: &crate::json::JsonValue,
ordinality: usize,
vars: Option<&crate::json::JsonValue>,
out: &mut alloc::vec::Vec<Row<'static>>,
) -> Result<(), EngineError> {
use spg_sql::ast::JsonTableColumn as C;
let mut parent_cells: alloc::vec::Vec<Value<'static>> = alloc::vec::Vec::new();
let mut nested_runs: alloc::vec::Vec<alloc::vec::Vec<Row<'static>>> =
alloc::vec::Vec::new();
let mut nested_widths: alloc::vec::Vec<usize> = alloc::vec::Vec::new();
for col in &jt.columns {
match col {
C::Ordinality { .. } => {
parent_cells.push(Value::BigInt(ordinality as i64 + 1));
}
C::Regular { .. } => {
parent_cells.push(self.json_table_column_value(col, item, vars)?);
}
C::Nested { path, columns } => {
let sub = spg_sql::ast::JsonTable {
doc: jt.doc.clone(), row_path: path.clone(),
columns: columns.clone(),
passing: alloc::vec::Vec::new(),
};
let (nrows, nschema) = self.json_table_rows(&sub, Some(item))?;
nested_widths.push(nschema.len());
nested_runs.push(nrows);
}
}
}
if nested_runs.is_empty() {
out.push(Row::new(parent_cells));
return Ok(());
}
let before = out.len();
for (s_idx, run) in nested_runs.iter().enumerate() {
for nrow in run {
let mut cells = parent_cells.clone();
for (o_idx, w) in nested_widths.iter().enumerate() {
if o_idx == s_idx {
cells.extend(nrow.values.iter().cloned());
} else {
for _ in 0..*w {
cells.push(Value::Null);
}
}
}
out.push(Row::new(cells));
}
}
if out.len() == before {
let mut cells = parent_cells.clone();
for w in &nested_widths {
for _ in 0..*w {
cells.push(Value::Null);
}
}
out.push(Row::new(cells));
}
Ok(())
}
fn json_table_column_value(
&self,
col: &spg_sql::ast::JsonTableColumn,
item: &crate::json::JsonValue,
vars: Option<&crate::json::JsonValue>,
) -> Result<Value<'static>, EngineError> {
use spg_sql::ast::{JsonTableColumn as C, JsonTableOnBehavior as B};
let C::Regular {
name,
ty,
path,
exists,
format_json,
wrapper,
on_empty,
on_error,
} = col
else {
unreachable!("caller guards Regular");
};
let matches = crate::json::json_table_path(item, path, vars).map_err(EngineError::Eval)?;
if *exists {
return Ok(Value::Bool(!matches.is_empty()));
}
let empty_schema: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let ctx = EvalContext::new(&empty_schema, None);
let dummy = Row::new(alloc::vec::Vec::new());
let default_of = |b: &B| -> Result<Option<Value<'static>>, EngineError> {
match b {
B::Null => Ok(Some(Value::Null)),
B::Error => Ok(None),
B::Default(e) => Ok(Some(
eval::eval_expr(e, &dummy, &ctx).map_err(EngineError::Eval)?,
)),
}
};
if matches.is_empty() {
return match default_of(on_empty)? {
Some(v) => coerce_json_table_default(v, *ty, name),
None => Err(EngineError::Unsupported(alloc::format!(
"no SQL/JSON item found for JSON_TABLE column {name:?}"
))),
};
}
let first = &matches[0];
if *format_json {
let text = if *wrapper {
crate::json::JsonValue::Array(matches.clone()).canonical_json_text()
} else {
first.canonical_json_text()
};
return Ok(Value::Json(alloc::borrow::Cow::Owned(text)));
}
if first.is_json_null() {
return Ok(Value::Null);
}
let dt = crate::conversions::column_type_to_data_type(*ty);
let scalar = Value::Text(alloc::borrow::Cow::Owned(first.scalar_text()));
match crate::conversions::coerce_value(scalar, dt, name, 0) {
Ok(v) => Ok(v),
Err(e) => match default_of(on_error)? {
Some(v) => coerce_json_table_default(v, *ty, name),
None => Err(e),
},
}
}
pub(crate) fn table_fn_rows(
&self,
primary: &TableRef,
) -> Result<(alloc::vec::Vec<Row<'static>>, alloc::vec::Vec<ColumnSchema>), EngineError> {
let (fn_name, args) = primary
.table_fn_call
.as_deref()
.expect("caller guards table_fn_call.is_some()");
let empty_schema: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let ctx = EvalContext::new(&empty_schema, None);
let dummy_row = Row::new(alloc::vec::Vec::new());
let arg0: Option<Value<'static>> = match args.first() {
Some(e) => Some(eval::eval_expr(e, &dummy_row, &ctx).map_err(EngineError::Eval)?),
None => None,
};
match fn_name.as_str() {
"jsonb_populate_record"
| "json_populate_record"
| "jsonb_populate_recordset"
| "json_populate_recordset" => {
let type_name = match args.first() {
Some(Expr::Cast {
target: spg_sql::ast::CastTarget::Named(n),
..
}) => n.clone(),
_ => {
return Err(EngineError::Unsupported(alloc::format!(
"{fn_name}(): first argument must name a row type, \
e.g. NULL::mytable"
)));
}
};
let cat = self.active_catalog();
let cols: alloc::vec::Vec<ColumnSchema> = if let Some(t) = cat.get(&type_name) {
t.schema().columns.clone()
} else if let Some(c) = cat.composite_types().get(&type_name) {
c.fields
.iter()
.map(|(n, ty)| ColumnSchema::new(n.clone(), *ty, true))
.collect()
} else {
return Err(EngineError::Unsupported(alloc::format!(
"type \"{type_name}\" does not exist"
)));
};
let json_arg = match args.get(1) {
Some(e) => eval::eval_expr(e, &dummy_row, &ctx).map_err(EngineError::Eval)?,
None => Value::Null,
};
let docs: alloc::vec::Vec<Value<'static>> = if fn_name.ends_with("recordset") {
crate::json::array_element_rows(&json_arg, false, fn_name)
.map_err(EngineError::Eval)?
.into_iter()
.map(|s| s.map_or(Value::Null, Value::json))
.collect()
} else if matches!(json_arg, Value::Null) {
alloc::vec::Vec::new()
} else {
alloc::vec![json_arg]
};
let mut rows = alloc::vec::Vec::with_capacity(docs.len());
for doc in &docs {
let mut vals = alloc::vec::Vec::with_capacity(cols.len());
for c in &cols {
let raw = crate::json::path_get(doc, &Value::text(c.name.clone()), true)
.map_err(EngineError::Eval)?;
let v = if matches!(raw, Value::Null) {
Value::Null
} else {
crate::conversions::coerce_value(raw, c.ty, "", 0)
.map_err(|e| EngineError::Unsupported(alloc::format!("{e:?}")))?
};
vals.push(v);
}
rows.push(Row::new(vals));
}
Ok((rows, cols))
}
"pg_partition_tree" => {
let cols = alloc::vec![
ColumnSchema::new("relid".to_string(), DataType::Text, true),
ColumnSchema::new("parentrelid".to_string(), DataType::Text, true),
ColumnSchema::new("isleaf".to_string(), DataType::Bool, true),
ColumnSchema::new("level".to_string(), DataType::Int, true),
];
let Some(Value::Text(name)) = &arg0 else {
return Ok((alloc::vec::Vec::new(), cols));
};
let entries = crate::partition_walks::tree_of(self.active_catalog(), name.as_ref());
if entries.is_empty() && self.active_catalog().get(name.as_ref()).is_none() {
return Err(EngineError::Unsupported(alloc::format!(
"relation \"{name}\" does not exist"
)));
}
let rows = entries
.into_iter()
.map(|(relid, parent, isleaf, level)| {
Row::new(alloc::vec![
Value::text(relid),
parent.map_or(Value::Null, Value::text),
Value::Bool(isleaf),
#[allow(clippy::cast_possible_truncation)]
Value::Int(level as i32),
])
})
.collect();
Ok((rows, cols))
}
"pg_partition_ancestors" => {
let cols =
alloc::vec![ColumnSchema::new("relid".to_string(), DataType::Text, true)];
let Some(Value::Text(name)) = &arg0 else {
return Ok((alloc::vec::Vec::new(), cols));
};
let cat = self.active_catalog();
if cat.get(name.as_ref()).is_none() {
return Err(EngineError::Unsupported(alloc::format!(
"relation \"{name}\" does not exist"
)));
}
let in_tree = cat
.get(name.as_ref())
.is_some_and(|t| t.schema().partition_role.is_some());
let rows = if in_tree {
crate::partition_walks::ancestors_of(cat, name.as_ref())
.into_iter()
.map(|n| Row::new(alloc::vec![Value::text(n)]))
.collect()
} else {
alloc::vec::Vec::new()
};
Ok((rows, cols))
}
"ts_debug" => {
use crate::fts::{TokenType, TsDict};
let cols = alloc::vec![
ColumnSchema::new("alias".to_string(), DataType::Text, false),
ColumnSchema::new("description".to_string(), DataType::Text, false),
ColumnSchema::new("token".to_string(), DataType::Text, false),
ColumnSchema::new("dictionaries".to_string(), DataType::TextArray, false),
ColumnSchema::new("dictionary".to_string(), DataType::Text, true),
ColumnSchema::new("lexemes".to_string(), DataType::TextArray, true),
];
let (cfg_name, text) = match (&arg0, args.get(1)) {
(Some(Value::Text(c)), Some(t)) => {
let v = eval::eval_expr(t, &dummy_row, &ctx).map_err(EngineError::Eval)?;
(c.to_string(), crate::eval::value_to_text(&v))
}
(Some(v), None) => (
alloc::string::String::from("english"),
crate::eval::value_to_text(v),
),
_ => return Ok((alloc::vec::Vec::new(), cols)),
};
let english = match cfg_name
.trim()
.trim_start_matches("pg_catalog.")
.to_ascii_lowercase()
.as_str()
{
"english" => true,
"simple" => false,
other => {
return Err(EngineError::Unsupported(alloc::format!(
"text search configuration \"{other}\" does not exist"
)));
}
};
let rows = crate::fts::tokenize_typed(&text)
.into_iter()
.map(|tok| {
let dict = tok.ty.dictionary(english);
let dname = dict.map(|d| match d {
TsDict::Simple => "simple",
TsDict::EnglishStem => "english_stem",
});
let folded = tok.text.to_lowercase();
let lexemes = dict.map(|d| match d {
TsDict::Simple => alloc::vec![Some(folded.clone())],
TsDict::EnglishStem => {
if crate::fts::is_english_stopword(&folded) {
alloc::vec::Vec::new()
} else {
alloc::vec![Some(crate::fts::porter_stem(&folded))]
}
}
});
Row::new(alloc::vec![
Value::text(tok.ty.alias()),
Value::text(tok.ty.description()),
Value::text(tok.text),
Value::TextArray(
dname
.map(|n| alloc::vec![Some(alloc::string::String::from(n))])
.unwrap_or_default(),
),
dname.map_or(Value::Null, Value::text),
lexemes.map_or(Value::Null, Value::TextArray),
])
})
.collect();
let _ = TokenType::AsciiWord;
Ok((rows, cols))
}
"ts_token_type" => {
use crate::fts::TokenType as T;
let cols = alloc::vec![
ColumnSchema::new("tokid".to_string(), DataType::Int, false),
ColumnSchema::new("alias".to_string(), DataType::Text, false),
ColumnSchema::new("description".to_string(), DataType::Text, false),
];
if let Some(Value::Text(p)) = &arg0
&& !p.eq_ignore_ascii_case("default")
&& !p.eq_ignore_ascii_case("pg_catalog.default")
{
return Err(EngineError::Unsupported(alloc::format!(
"text search parser \"{p}\" does not exist"
)));
}
const TYPES: &[T] = &[
T::AsciiWord,
T::Word,
T::NumWord,
T::Email,
T::Url,
T::Host,
T::SFloat,
T::Version,
T::HwordNumPart,
T::HwordPart,
T::HwordAsciiPart,
T::Blank,
T::Tag,
T::Protocol,
T::NumHword,
T::AsciiHword,
T::Hword,
T::UrlPath,
T::File,
T::Float,
T::Int,
T::Uint,
T::Entity,
];
let rows = TYPES
.iter()
.map(|t| {
Row::new(alloc::vec![
Value::Int(*t as i32),
Value::text(t.alias()),
Value::text(t.description()),
])
})
.collect();
Ok((rows, cols))
}
other => {
if !self.active_catalog().functions_named(other).is_empty() {
return self.exec_setof_user_function(other, args, primary.alias.as_deref());
}
Err(EngineError::Unsupported(alloc::format!(
"table function {other}() is not supported in FROM"
)))
}
}
}
fn exec_setof_user_function(
&self,
name: &str,
args: &[spg_sql::ast::Expr],
alias: Option<&str>,
) -> Result<(alloc::vec::Vec<Row<'static>>, alloc::vec::Vec<ColumnSchema>), EngineError> {
let empty: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let arg_ctx = self.ev_ctx(&empty, None);
let dummy = Row::new(alloc::vec::Vec::new());
let mut vals: alloc::vec::Vec<Value<'static>> = alloc::vec::Vec::new();
for a in args {
vals.push(eval::eval_expr(a, &dummy, &arg_ctx).map_err(EngineError::Eval)?);
}
self.setof_rows_of(name, &vals, alias)
}
pub(crate) fn setof_rows_of(
&self,
name: &str,
arg_values: &[Value<'static>],
alias: Option<&str>,
) -> Result<(alloc::vec::Vec<Row<'static>>, alloc::vec::Vec<ColumnSchema>), EngineError> {
let cat = self.active_catalog();
let overloads = cat.functions_named(name);
let def = overloads
.iter()
.find(|f| spg_storage::function_arg_types(&f.args_repr).len() == arg_values.len())
.ok_or_else(|| {
EngineError::Unsupported(alloc::format!(
"function {name} does not exist with {} argument(s)",
arg_values.len()
))
})?;
let declared = def.returns.trim().to_string();
let upper = declared.to_ascii_uppercase();
if !upper.starts_with("SETOF") && !upper.starts_with("TABLE(") {
return Err(EngineError::Unsupported(alloc::format!(
"function {name}() does not return a set — it cannot be used in FROM"
)));
}
let arg_names_pl = spg_storage::function_arg_names(&def.args_repr);
if def.language.eq_ignore_ascii_case("plpgsql") {
let out_rows = self
.call_plpgsql_setof_fn(def, &arg_names_pl, arg_values)
.map_err(EngineError::Eval)?;
let cols = setof_column_shape(&declared, name, alias, out_rows.first());
let rows = out_rows.into_iter().map(Row::new).collect();
return Ok((rows, cols));
}
let body = def.body.trim().trim_end_matches(';');
let stmt = spg_sql::parser::parse_statement(body).map_err(|e| {
EngineError::Unsupported(alloc::format!("function {name} body does not parse: {e}"))
})?;
let spg_sql::ast::Statement::Select(body_select) = stmt else {
return Err(EngineError::Unsupported(alloc::format!(
"function {name}(): a set-returning body must be a SELECT"
)));
};
let arg_names = spg_storage::function_arg_names(&def.args_repr);
let bound = crate::eval::bind_user_fn_args(
self.active_catalog(),
&body_select,
&arg_names,
arg_values,
)
.map_err(EngineError::Eval)?;
let out = self.exec_select_cancel(&bound, crate::CancelToken::none())?;
let QueryResult::Rows { columns, rows } = out else {
return Ok((alloc::vec::Vec::new(), alloc::vec::Vec::new()));
};
let cols = setof_column_shape_from(&declared, name, alias, &columns);
Ok((rows, cols))
}
fn exec_select_jsonb_each_text(
&self,
stmt: &SelectStatement,
primary: &TableRef,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
let (each_fn, arg_expr) = primary
.jsonb_each_text_arg
.as_ref()
.map(|(name, expr)| (name.as_str(), expr.as_ref()))
.expect("caller guards jsonb_each_text_arg.is_some()");
let as_text = each_fn.ends_with("_text");
let empty_schema: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let ctx = EvalContext::new(&empty_schema, None);
let dummy_row = Row::new(alloc::vec::Vec::new());
let arg_value = eval::eval_expr(arg_expr, &dummy_row, &ctx).map_err(EngineError::Eval)?;
let pairs =
crate::json::each_rows(&arg_value, as_text, each_fn).map_err(EngineError::Eval)?;
let rows: alloc::vec::Vec<Row<'static>> = pairs
.into_iter()
.map(|(k, v)| {
let key_val = Value::text(k);
let value_val = match v {
Some(s) if as_text => Value::text(s),
Some(s) => Value::Json(alloc::borrow::Cow::Owned(s)),
None => Value::Null,
};
Row::new(alloc::vec![key_val, value_val])
})
.collect();
let alias = primary.alias.clone().unwrap_or_else(|| each_fn.to_string());
let value_dtype = if as_text {
spg_storage::DataType::Text
} else {
spg_storage::DataType::Json
};
let key_col = ColumnSchema::new("key".to_string(), spg_storage::DataType::Text, false);
let value_col = ColumnSchema::new("value".to_string(), value_dtype, as_text);
let mut schema_cols = alloc::vec![key_col, value_col];
for (i, new_name) in primary.unnest_column_aliases.iter().enumerate() {
if let Some(col) = schema_cols.get_mut(i) {
col.name = new_name.clone();
}
}
let scan_ctx = self.ev_ctx(&schema_cols, Some(&alias));
let filtered: alloc::vec::Vec<Row<'static>> = if let Some(w) = &stmt.where_ {
let mut out = alloc::vec::Vec::with_capacity(rows.len());
for row in rows {
cancel.check()?;
let v = eval::eval_expr(w, &row, &scan_ctx).map_err(EngineError::Eval)?;
if matches!(v, Value::Bool(true)) {
out.push(row);
}
}
out
} else {
rows
};
if aggregate::uses_aggregate(stmt) {
let agg_memo = core::cell::RefCell::new(memoize::MemoizeCache::default());
let agg_correlated = |e: &Expr, r: &Row<'static>, c: &EvalContext<'_>| {
self.eval_expr_with_correlated(e, r, c, cancel, Some(&mut agg_memo.borrow_mut()))
.map_err(|err| match err {
EngineError::Eval(ev) => ev,
other => eval::EvalError::TypeMismatch {
detail: alloc::format!("{other}"),
},
})
};
let agg = aggregate::run(
stmt,
crate::join::AggRows::Owned(&filtered),
&schema_cols,
Some(&alias),
Some(&agg_correlated),
self.parallel_runner.0.as_deref(),
Some(self.active_catalog()),
Some(self),
)?;
return self.finish_agg_result(agg, stmt, cancel);
}
let projection =
build_projection(&stmt.items, &schema_cols, &alias, self.backslash_escapes)?;
let mut projected_rows: alloc::vec::Vec<Row<'static>> =
alloc::vec::Vec::with_capacity(filtered.len());
for row in &filtered {
let mut vals = alloc::vec::Vec::with_capacity(projection.len());
for p in &projection {
let v = eval::eval_expr(&p.expr, row, &scan_ctx).map_err(EngineError::Eval)?;
vals.push(v);
}
projected_rows.push(Row::new(vals));
}
let columns: alloc::vec::Vec<ColumnSchema> = projection
.iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name.clone(), p.ty, p.nullable);
c.user_enum_type = p.user_enum_type.clone();
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
if !stmt.order_by.is_empty() {
let mut indexed: alloc::vec::Vec<(usize, Vec<Value<'static>>)> = filtered
.iter()
.enumerate()
.map(|(i, r)| -> Result<_, EngineError> {
let keys: Result<Vec<Value<'static>>, EngineError> = stmt
.order_by
.iter()
.map(|ob| {
eval::eval_expr(&ob.expr, r, &scan_ctx).map_err(EngineError::Eval)
})
.collect();
Ok((i, keys?))
})
.collect::<Result<_, _>>()?;
indexed.sort_by(|a, b| {
for (idx, (ka, kb)) in a.1.iter().zip(b.1.iter()).enumerate() {
let o = &stmt.order_by[idx];
let cmp = order_by_value_cmp_in(
o.desc,
o.nulls_first,
ka,
kb,
scan_ctx.mysql_dialect && !crate::eval::is_binary_coerced(&o.expr),
);
if cmp != core::cmp::Ordering::Equal {
return cmp;
}
}
core::cmp::Ordering::Equal
});
projected_rows = indexed
.into_iter()
.map(|(i, _)| projected_rows[i].clone())
.collect();
}
if stmt.distinct {
projected_rows = dedup_rows(projected_rows, scan_ctx.mysql_dialect);
}
if let Some(offset) = stmt.offset_literal() {
let off = (offset as usize).min(projected_rows.len());
projected_rows.drain(..off);
}
if let Some(limit) = stmt.limit_literal() {
projected_rows.truncate(limit as usize);
}
Ok(QueryResult::Rows {
columns,
rows: projected_rows,
})
}
fn exec_select_derived(
&self,
stmt: &SelectStatement,
primary: &TableRef,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
let inner = primary
.lateral_subquery
.as_deref()
.expect("caller guards lateral_subquery.is_some()");
let QueryResult::Rows {
columns: inner_cols,
rows,
} = self.exec_select_cancel(inner, cancel)?
else {
return Err(EngineError::Unsupported(
"derived table subquery must return rows".into(),
));
};
let alias = primary
.alias
.clone()
.unwrap_or_else(|| primary.name.clone());
let mut schema_cols: alloc::vec::Vec<ColumnSchema> = inner_cols;
let n_out = schema_cols.len() + usize::from(primary.with_ordinality);
if primary.unnest_column_aliases.len() > n_out {
return Err(EngineError::Unsupported(alloc::format!(
"table \"{alias}\" has {n_out} columns available but {} columns specified",
primary.unnest_column_aliases.len()
)));
}
if primary.scalar_fn_item && schema_cols.len() == 1 {
schema_cols[0].scalar_row_source = true;
}
let mut rows = rows;
if primary.with_ordinality {
schema_cols.push(ColumnSchema::new(
"ordinality".to_string(),
DataType::BigInt,
false,
));
rows = rows
.into_iter()
.enumerate()
.map(|(i, r)| {
let mut v = r.values;
#[allow(clippy::cast_possible_wrap)]
v.push(Value::BigInt(i as i64 + 1));
Row::new(v)
})
.collect();
}
for (i, new_name) in primary.unnest_column_aliases.iter().enumerate() {
if let Some(col) = schema_cols.get_mut(i) {
col.name = new_name.clone();
}
}
self.exec_select_over_rows(stmt, rows, schema_cols, &alias, cancel)
}
fn exec_select_over_rows(
&self,
stmt: &SelectStatement,
rows: alloc::vec::Vec<Row<'static>>,
schema_cols: alloc::vec::Vec<ColumnSchema>,
alias: &str,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
let scan_ctx = self.ev_ctx(&schema_cols, Some(alias));
let corr_memo = core::cell::RefCell::new(memoize::MemoizeCache::default());
let filtered: alloc::vec::Vec<Row<'static>> = if let Some(w) = &stmt.where_ {
let mut out = alloc::vec::Vec::with_capacity(rows.len());
for row in rows {
cancel.check()?;
let v = self.eval_expr_with_correlated(
w,
&row,
&scan_ctx,
cancel,
Some(&mut corr_memo.borrow_mut()),
)?;
if matches!(v, Value::Bool(true)) {
out.push(row);
}
}
out
} else {
rows
};
if aggregate::uses_aggregate(stmt) {
let agg_memo = core::cell::RefCell::new(memoize::MemoizeCache::default());
let agg_correlated = |e: &Expr, r: &Row<'static>, c: &EvalContext<'_>| {
self.eval_expr_with_correlated(e, r, c, cancel, Some(&mut agg_memo.borrow_mut()))
.map_err(|err| match err {
EngineError::Eval(ev) => ev,
other => eval::EvalError::TypeMismatch {
detail: alloc::format!("{other}"),
},
})
};
let agg = aggregate::run(
stmt,
crate::join::AggRows::Owned(&filtered),
&schema_cols,
Some(alias),
Some(&agg_correlated),
self.parallel_runner.0.as_deref(),
Some(self.active_catalog()),
Some(self),
)?;
return self.finish_agg_result(agg, stmt, cancel);
}
let projection =
build_projection(&stmt.items, &schema_cols, alias, self.backslash_escapes)?;
let srf_idxs = self.srf_target_idxs(&projection);
let mut src_of_row: alloc::vec::Vec<usize> = alloc::vec::Vec::new();
let mut projected_rows: alloc::vec::Vec<Row<'static>> =
alloc::vec::Vec::with_capacity(filtered.len());
if !srf_idxs.is_empty() {
let (rows, src) =
expand_projection_srfs(self, &projection, &srf_idxs, &filtered, &scan_ctx)?;
projected_rows = rows;
src_of_row = src;
} else {
for row in &filtered {
let mut vals = alloc::vec::Vec::with_capacity(projection.len());
for p in &projection {
let v = self.eval_expr_with_correlated(
&p.expr,
row,
&scan_ctx,
cancel,
Some(&mut corr_memo.borrow_mut()),
)?;
vals.push(v);
}
projected_rows.push(Row::new(vals));
}
}
let columns: alloc::vec::Vec<ColumnSchema> = projection
.iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name.clone(), p.ty, p.nullable);
c.user_enum_type = p.user_enum_type.clone();
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
let order_by = resolve_positional_order_by(&stmt.order_by, &projection);
if !order_by.is_empty() {
let out_cols = if srf_idxs.is_empty() {
alloc::vec![None; order_by.len()]
} else {
srf_order_output_cols(&order_by, &projection)
};
let mut indexed: alloc::vec::Vec<(usize, Vec<Value<'static>>)> = projected_rows
.iter()
.enumerate()
.map(|(k, out)| -> Result<_, EngineError> {
let r = &filtered[src_of_row.get(k).copied().unwrap_or(k)];
let keys: Result<Vec<Value<'static>>, EngineError> = order_by
.iter()
.zip(out_cols.iter())
.map(|(ob, oc)| {
let v = srf_order_key(ob, *oc, out, r, &scan_ctx)?;
Ok(
match crate::orderby::enum_order_ordinal(&ob.expr, &v, &scan_ctx) {
Some(ord) => Value::Float(ord),
None => v,
},
)
})
.collect();
Ok((k, keys?))
})
.collect::<Result<_, _>>()?;
indexed.sort_by(|a, b| {
for (idx, (ka, kb)) in a.1.iter().zip(b.1.iter()).enumerate() {
let o = &stmt.order_by[idx];
let cmp = order_by_value_cmp_in(
o.desc,
o.nulls_first,
ka,
kb,
scan_ctx.mysql_dialect && !crate::eval::is_binary_coerced(&o.expr),
);
if cmp != core::cmp::Ordering::Equal {
return cmp;
}
}
core::cmp::Ordering::Equal
});
projected_rows = indexed
.into_iter()
.map(|(i, _)| projected_rows[i].clone())
.collect();
}
if stmt.distinct {
projected_rows = dedup_rows(projected_rows, scan_ctx.mysql_dialect);
}
if let Some(offset) = stmt.offset_literal() {
let off = (offset as usize).min(projected_rows.len());
projected_rows.drain(..off);
}
if let Some(limit) = stmt.limit_literal() {
projected_rows.truncate(limit as usize);
}
Ok(QueryResult::Rows {
columns,
rows: projected_rows,
})
}
fn exec_constant_select(&self, stmt: &SelectStatement) -> Result<QueryResult, EngineError> {
let empty_schema: Vec<ColumnSchema> = Vec::new();
let ctx = self.ev_ctx(&empty_schema, None);
if aggregate::uses_aggregate(stmt) {
let dummy = Row::new(Vec::new());
let passes = match &stmt.where_ {
Some(w) => matches!(eval::eval_expr(w, &dummy, &ctx)?, Value::Bool(true)),
None => true,
};
let rows: Vec<RowRef<'_>> = if passes {
alloc::vec![RowRef::Owned(&dummy)]
} else {
Vec::new()
};
let agg = aggregate::run(
stmt,
crate::join::AggRows::Refs(&rows),
&empty_schema,
None,
None,
self.parallel_runner.0.as_deref(),
Some(self.active_catalog()),
Some(self),
)?;
return self.finish_agg_result(agg, stmt, CancelToken::none());
}
let projection = build_projection(&stmt.items, &empty_schema, "", self.backslash_escapes)?;
let dummy_row = Row::new(Vec::new());
if let Some(w) = &stmt.where_ {
let cond = eval::eval_expr(w, &dummy_row, &ctx)?;
if !crate::eval::predicate_is_true(&cond, "WHERE", ctx.mysql_dialect)? {
let columns: Vec<ColumnSchema> = projection
.into_iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name, p.ty, p.nullable);
c.user_enum_type = p.user_enum_type;
c.collation_name = p.collation_name;
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
return Ok(QueryResult::Rows {
columns,
rows: Vec::new(),
});
}
}
let srf_idxs = self.srf_target_idxs(&projection);
if !srf_idxs.is_empty() {
let mut rows = expand_srf_row(self, &projection, &srf_idxs, &dummy_row, &ctx)?;
let columns: Vec<ColumnSchema> = projection
.into_iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name, p.ty, p.nullable);
c.user_enum_type = p.user_enum_type;
c.collation_name = p.collation_name;
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
if !stmt.order_by.is_empty() {
let synth_ctx =
EvalContext::new(&columns, None).with_catalog(self.active_catalog());
let resolved: Vec<spg_sql::ast::OrderBy> = stmt
.order_by
.iter()
.map(|o| {
let mut o = o.clone();
if let Expr::Literal(spg_sql::ast::Literal::Integer(n)) = &o.expr
&& *n >= 1
&& let Ok(idx) = usize::try_from(*n - 1)
&& idx < columns.len()
{
o.expr = Expr::Column(spg_sql::ast::ColumnName {
qualifier: None,
name: columns[idx].name.clone(),
});
}
o
})
.collect();
let descs: Vec<bool> = resolved.iter().map(|o| o.desc).collect();
let mut tagged: Vec<(Vec<OrderKey>, Row)> = Vec::with_capacity(rows.len());
for r in rows {
let keys = build_order_keys(&resolved, &r, &synth_ctx)?;
tagged.push((keys, r));
}
sort_by_keys(&mut tagged, &descs);
rows = tagged.into_iter().map(|(_, r)| r).collect();
}
apply_offset_and_limit(&mut rows, stmt.offset_literal(), stmt.limit_literal());
return Ok(QueryResult::Rows { columns, rows });
}
let mut values = Vec::with_capacity(projection.len());
for p in &projection {
values.push(eval::eval_expr(&p.expr, &dummy_row, &ctx)?);
}
let columns: Vec<ColumnSchema> = projection
.into_iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name, p.ty, p.nullable);
c.user_enum_type = p.user_enum_type;
c.collation_name = p.collation_name;
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
let mut rows = alloc::vec![Row::new(values)];
apply_offset_and_limit(&mut rows, stmt.offset_literal(), stmt.limit_literal());
Ok(QueryResult::Rows { columns, rows })
}
pub(crate) fn try_count_star_pk_in_subquery_fast(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<Option<QueryResult>, EngineError> {
use spg_sql::ast::SelectItem;
if stmt.distinct
|| stmt.limit_with_ties
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| !stmt.unions.is_empty()
|| !stmt.order_by.is_empty()
|| stmt.limit.is_some()
|| stmt.offset.is_some()
|| stmt.items.len() != 1
{
return Ok(None);
}
let SelectItem::Expr { expr, .. } = &stmt.items[0] else {
return Ok(None);
};
let is_count_star = matches!(expr, Expr::FunctionCall { name, args }
if name.eq_ignore_ascii_case("count_star") && args.is_empty());
if !is_count_star {
return Ok(None);
}
let Some(from) = stmt.from.as_ref() else {
return Ok(None);
};
if !from.joins.is_empty()
|| from.primary.lateral_subquery.is_some()
|| from.primary.unnest_expr.is_some()
|| from.primary.generate_series_args.is_some()
|| from.primary.table_fn_call.is_some()
|| from.primary.as_of_segment.is_some()
{
return Ok(None);
}
let Some(where_expr) = stmt.where_.as_ref() else {
return Ok(None);
};
let Expr::InSubquery {
expr: col_expr,
subquery,
negated: false,
} = where_expr
else {
return Ok(None);
};
let Expr::Column(c) = col_expr.as_ref() else {
return Ok(None);
};
let outer_alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str());
if let Some(q) = c.qualifier.as_deref()
&& !q.eq_ignore_ascii_case(outer_alias)
{
return Ok(None);
}
let catalog = self.active_catalog();
let Some(outer_table) = catalog.get(from.primary.name.as_str()) else {
return Ok(None);
};
let outer_schema = outer_table.schema();
let Some(outer_pos) = outer_schema
.columns
.iter()
.position(|s| s.name.eq_ignore_ascii_case(&c.name))
else {
return Ok(None);
};
if !matches!(
outer_schema.columns[outer_pos].ty,
spg_storage::DataType::BigInt
| spg_storage::DataType::Int
| spg_storage::DataType::SmallInt
) {
return Ok(None);
}
if !outer_schema
.uniqueness_constraints
.iter()
.any(|u| u.is_primary_key && u.columns.as_slice() == [outer_pos])
{
return Ok(None);
}
let Some(idx) = outer_table.index_on(outer_pos) else {
return Ok(None);
};
if crate::subquery::select_is_correlated(subquery) {
return Ok(None);
}
let mut inner = (**subquery).clone();
self.resolve_select_subqueries(&mut inner, cancel)?;
let r = match self.exec_bare_select_cancel(&inner, cancel) {
Ok(r) => r,
Err(_) => return Ok(None),
};
let QueryResult::Rows { columns, rows, .. } = r else {
return Ok(None);
};
if columns.len() != 1 {
return Ok(None);
}
let inner_unique = (|| -> bool {
if inner.distinct
|| inner.group_by.is_some()
|| !inner.unions.is_empty()
|| inner.having.is_some()
|| inner.items.len() != 1
{
return false;
}
let Some(inner_from) = inner.from.as_ref() else {
return false;
};
if !inner_from.joins.is_empty()
|| inner_from.primary.lateral_subquery.is_some()
|| inner_from.primary.unnest_expr.is_some()
|| inner_from.primary.generate_series_args.is_some()
|| inner_from.primary.table_fn_call.is_some()
{
return false;
}
let SelectItem::Expr { expr: proj, .. } = &inner.items[0] else {
return false;
};
let Expr::Column(pc) = proj else {
return false;
};
let inner_alias = inner_from
.primary
.alias
.as_deref()
.unwrap_or(inner_from.primary.name.as_str());
if let Some(q) = pc.qualifier.as_deref()
&& !q.eq_ignore_ascii_case(inner_alias)
{
return false;
}
let Some(inner_table) = catalog.get(inner_from.primary.name.as_str()) else {
return false;
};
let isch = inner_table.schema();
let Some(ipos) = isch
.columns
.iter()
.position(|s| s.name.eq_ignore_ascii_case(&pc.name))
else {
return false;
};
isch.uniqueness_constraints
.iter()
.any(|u| u.columns.as_slice() == [ipos])
})();
let mut count: i64 = 0;
let mut probed = if inner_unique {
hashbrown::HashSet::<i64>::new()
} else {
hashbrown::HashSet::<i64>::with_capacity(rows.len())
};
for row in &rows {
let v = row.values.first().cloned().unwrap_or(Value::Null);
let n = match v {
Value::BigInt(n) => n,
Value::Int(n) => i64::from(n),
Value::SmallInt(n) => i64::from(n),
Value::Null => continue,
_ => return Ok(None),
};
if !inner_unique && !probed.insert(n) {
continue;
}
if !idx.lookup_eq_i64(n).is_empty() {
count += 1;
}
}
let columns_out = alloc::vec![ColumnSchema::new(
"count".to_string(),
spg_storage::DataType::BigInt,
false,
)];
let rows_out = alloc::vec![Row::new(alloc::vec![Value::BigInt(count)])];
Ok(Some(QueryResult::Rows {
columns: columns_out,
rows: rows_out,
}))
}
fn try_count_star_pk_in_list_fast(
&self,
stmt: &SelectStatement,
table: &spg_storage::Table,
schema_cols: &[ColumnSchema],
alias: &str,
) -> Option<QueryResult> {
use spg_sql::ast::{ColumnName, SelectItem};
if stmt.distinct
|| stmt.limit_with_ties
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| !stmt.unions.is_empty()
|| !stmt.order_by.is_empty()
|| stmt.limit.is_some()
|| stmt.offset.is_some()
|| stmt.items.len() != 1
{
return None;
}
let SelectItem::Expr { expr, .. } = &stmt.items[0] else {
return None;
};
let is_count_star = matches!(expr, Expr::FunctionCall { name, args }
if name.eq_ignore_ascii_case("count_star") && args.is_empty());
if !is_count_star {
return None;
}
let where_expr = stmt.where_.as_ref()?;
let Expr::InList {
expr: col_expr,
list,
negated: false,
} = where_expr
else {
return None;
};
let Expr::Column(c) = col_expr.as_ref() else {
return None;
};
if let Some(q) = c.qualifier.as_deref()
&& !q.eq_ignore_ascii_case(alias)
{
return None;
}
let col_pos = schema_cols
.iter()
.position(|s| s.name.eq_ignore_ascii_case(&c.name))?;
let schema = table.schema();
if !matches!(
schema.columns[col_pos].ty,
spg_storage::DataType::BigInt
| spg_storage::DataType::Int
| spg_storage::DataType::SmallInt
) {
return None;
}
if !schema
.uniqueness_constraints
.iter()
.any(|u| u.is_primary_key && u.columns.as_slice() == [col_pos])
{
return None;
}
let idx = table.index_on(col_pos)?;
let mut count: i64 = 0;
for lit in list {
let Expr::Literal(l) = lit else {
return None;
};
let v = eval::literal_to_value(l);
let key = spg_storage::IndexKey::from_value(&v)?;
if !idx.lookup_eq(&key).is_empty() {
count += 1;
}
}
let columns = alloc::vec![ColumnSchema::new(
"count".to_string(),
spg_storage::DataType::BigInt,
false,
)];
let rows = alloc::vec![Row::new(alloc::vec![Value::BigInt(count)])];
let _ = ColumnName {
qualifier: None,
name: String::new(),
};
Some(QueryResult::Rows { columns, rows })
}
fn try_count_star_indexed_range_fast(
&self,
stmt: &SelectStatement,
table: &spg_storage::Table,
schema_cols: &[ColumnSchema],
alias: &str,
snapshot: &spg_storage::snapshot::Snapshot,
) -> Option<QueryResult> {
use spg_sql::ast::SelectItem;
if stmt.distinct
|| stmt.limit_with_ties
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| !stmt.unions.is_empty()
|| !stmt.order_by.is_empty()
|| stmt.limit.is_some()
|| stmt.offset.is_some()
|| stmt.items.len() != 1
{
return None;
}
let SelectItem::Expr { expr, .. } = &stmt.items[0] else {
return None;
};
let is_count_star = matches!(expr, Expr::FunctionCall { name, args }
if name.eq_ignore_ascii_case("count_star") && args.is_empty());
if !is_count_star {
return None;
}
let where_expr = stmt.where_.as_ref()?;
let count =
crate::index_access::try_range_count(where_expr, schema_cols, table, alias, snapshot)?;
let columns = alloc::vec![ColumnSchema::new(
"count".to_string(),
spg_storage::DataType::BigInt,
false,
)];
let rows = alloc::vec![Row::new(alloc::vec![Value::BigInt(count)])];
Some(QueryResult::Rows { columns, rows })
}
fn run_single_table_aggregate<'a>(
&self,
stmt: &SelectStatement,
table: &'a spg_storage::Table,
schema_cols: &'a [ColumnSchema],
alias: &str,
indexed_rows: Option<Vec<Cow<'a, Row<'static>>>>,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
let sample_cell: core::cell::Cell<Option<u64>> = core::cell::Cell::new(None);
let ctx = self
.ev_ctx(schema_cols, Some(alias))
.with_sample_rng(&sample_cell);
let mut filtered: Vec<&Row<'static>> = if stmt.where_.is_none() {
Vec::with_capacity(table.rows().len())
} else {
Vec::new()
};
let mut memo = memoize::MemoizeCache::new();
let compiled_where: Option<eval::CompiledExpr> = stmt
.where_
.as_ref()
.filter(|w| eval::fully_compilable(w))
.map(|w| eval::compile_expr(w, &ctx));
let mut eval_stack: Vec<Value<'static>> = Vec::new();
let mut row_passes_where = |row: &Row<'static>,
eval_stack: &mut Vec<Value<'static>>,
memo: &mut memoize::MemoizeCache|
-> Result<bool, EngineError> {
match (&compiled_where, &stmt.where_) {
(Some(cw), _) => {
Ok(eval::compiled::eval_compiled_pred(
cw,
row,
&ctx,
eval_stack,
ctx.mysql_dialect,
)
.map_err(EngineError::Eval)?)
}
(None, Some(w)) => {
let cond = self.eval_expr_with_correlated(w, row, &ctx, cancel, Some(memo))?;
Ok(crate::eval::predicate_is_true(
&cond,
"WHERE",
ctx.mysql_dialect,
)?)
}
(None, None) => Ok(true),
}
};
if let Some(rows) = &indexed_rows {
for cow in rows {
let row = cow.as_ref();
if !row_passes_where(row, &mut eval_stack, &mut memo)? {
continue;
}
filtered.push(row);
}
}
let cold_rows_storage = if indexed_rows.is_none() {
self.iter_cold_rows_of_table(table)
} else {
Vec::new()
};
if indexed_rows.is_none() {
let scan_snapshot = self.current_snapshot();
table.note_seq_scan();
let n = table.row_count();
let par = self.parallel_runner.0.as_deref().filter(|_| {
n >= crate::PARALLEL_MIN_ROWS && (stmt.where_.is_none() || compiled_where.is_some())
});
if let Some(r) = par {
let n_shards = (n / crate::PARALLEL_MIN_ROWS).clamp(2, 8);
let chunk = n.div_ceil(n_shards);
type ShardOut = Result<alloc::vec::Vec<usize>, EngineError>;
let cw = &compiled_where;
let snap_ref = &scan_snapshot;
let results = r.run_shards(n_shards, &|s| {
let lo = s * chunk;
let hi = ((s + 1) * chunk).min(n);
let mut keep: alloc::vec::Vec<usize> = alloc::vec::Vec::with_capacity(hi - lo);
let shard_ctx = EvalContext::new(schema_cols, Some(alias));
let mut stack: Vec<Value<'static>> = Vec::new();
let out: ShardOut = (|| {
for i in lo..hi {
if !table.is_row_visible(i, snap_ref) {
continue;
}
let row = &table.rows()[i];
let pass = match cw {
Some(c) => eval::compiled::eval_compiled_pred(
c,
row,
&shard_ctx,
&mut stack,
shard_ctx.mysql_dialect,
)
.map_err(EngineError::Eval)?,
None => true,
};
if pass {
keep.push(i);
}
}
Ok(keep)
})();
alloc::boxed::Box::new(out)
});
let mut rows_cur = table.rows().run_cursor();
for boxed in results {
let shard = boxed
.downcast::<ShardOut>()
.expect("runner echoes the closure's box");
for i in (*shard)? {
if let Some(row) = rows_cur.get(i) {
filtered.push(row);
}
}
}
} else {
let mut rows_cur = table.rows().run_cursor();
for i in 0..n {
if !table.is_row_visible(i, &scan_snapshot) {
continue;
}
let Some(row) = rows_cur.get(i) else { continue };
if !row_passes_where(row, &mut eval_stack, &mut memo)? {
continue;
}
filtered.push(row);
}
}
for row in &cold_rows_storage {
if !row_passes_where(row, &mut eval_stack, &mut memo)? {
continue;
}
filtered.push(row);
}
}
let agg_memo = core::cell::RefCell::new(memoize::MemoizeCache::default());
let agg_correlated = |e: &Expr, r: &Row<'static>, c: &EvalContext<'_>| {
self.eval_expr_with_correlated(e, r, c, cancel, Some(&mut agg_memo.borrow_mut()))
.map_err(|err| match err {
EngineError::Eval(ev) => ev,
other => eval::EvalError::TypeMismatch {
detail: alloc::format!("{other}"),
},
})
};
let agg = aggregate::run(
stmt,
crate::join::AggRows::Ptrs(&filtered),
schema_cols,
Some(alias),
Some(&agg_correlated),
self.parallel_runner.0.as_deref(),
Some(self.active_catalog()),
Some(self),
)?;
self.finish_agg_result(agg, stmt, cancel)
}
fn run_single_table_scan<'a>(
&self,
stmt: &SelectStatement,
table: &'a spg_storage::Table,
schema_cols: &'a [ColumnSchema],
alias: &str,
indexed_rows: Option<Vec<Cow<'a, Row<'static>>>>,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
let sample_cell: core::cell::Cell<Option<u64>> = core::cell::Cell::new(None);
let ctx = self
.ev_ctx(schema_cols, Some(alias))
.with_sample_rng(&sample_cell);
let projection = build_projection(&stmt.items, schema_cols, alias, self.backslash_escapes)?;
let srf_idxs = self.srf_target_idxs(&projection);
let srf_position = srf_idxs.first().copied();
let mut srf_plan = if srf_position.is_some() {
Some(build_srf_plan(self, &projection, &srf_idxs, &ctx)?)
} else {
None
};
let mut tagged: Vec<(Vec<OrderKey>, Row<'static>)> = Vec::new();
let mut budget = ByteBudget::new(self.max_query_bytes);
let mut memo = memoize::MemoizeCache::new();
let compiled_where: Option<eval::CompiledExpr> = stmt
.where_
.as_ref()
.filter(|w| eval::fully_compilable(w))
.map(|w| eval::compile_expr(w, &ctx));
let mut eval_stack: Vec<Value<'static>> = Vec::new();
let scalarsq_fast: Vec<Option<crate::ScalarPkProbeFastPath>> = projection
.iter()
.map(|p| {
if let Expr::ScalarSubquery(inner) = &p.expr {
self.analyse_scalar_count_pk_eq_probe(inner, schema_cols, alias)
} else {
None
}
})
.collect();
let any_scalarsq_fast = scalarsq_fast.iter().any(Option::is_some);
let proj_direct = bind_direct_columns(&projection, &ctx);
let any_proj_direct = proj_direct.iter().any(Option::is_some);
let proj_const: Vec<Option<Value<'static>>> = projection
.iter()
.map(|p| crate::eval::compiled::constant_projection_value(&p.expr, &ctx))
.collect();
let any_proj_const = proj_const.iter().any(Option::is_some);
crate::bump_counter!(crate::select::SCAN_PATH_ENTERED);
let order_by = resolve_positional_order_by(&stmt.order_by, &projection);
let srf_order_cols: Vec<Option<usize>> = if srf_position.is_some() {
srf_order_output_cols(&order_by, &projection)
} else {
Vec::new()
};
let srf_key_bound: Vec<Option<usize>> = (0..order_by.len()).map(Some).collect();
let early_cap: Option<usize> = if order_by.is_empty()
&& !stmt.distinct
&& !stmt.limit_with_ties
&& srf_position.is_none()
&& stmt.where_.is_none()
{
stmt.limit_literal()
.map(|n| n.saturating_add(stmt.offset_literal().unwrap_or(0)) as usize)
} else {
None
};
let order_colls = crate::orderby::order_by_collations(&order_by, &ctx)?;
let topk_stream: Option<(usize, Vec<bool>)> = if !order_by.is_empty()
&& !stmt.distinct
&& !stmt.limit_with_ties
&& srf_position.is_none()
&& !self.env_cfg().disable_topk
{
stmt.limit_literal().and_then(|l| {
let keep = (l as usize).saturating_add(stmt.offset_literal().unwrap_or(0) as usize);
(keep >= 1).then(|| (keep, order_by.iter().map(|o| o.desc).collect()))
})
} else {
None
};
let mut seen_distinct: hashbrown::HashMap<u64, alloc::vec::Vec<usize>> =
hashbrown::HashMap::new();
let distinct_hb = hashbrown::DefaultHashBuilder::default();
let mut proj_buf: Vec<Value<'static>> = Vec::new();
let mut proj_pool: Vec<Vec<Value<'static>>> = Vec::new();
let mut key_pool: Vec<Vec<crate::orderby::OrderKey>> = Vec::new();
let mut topk_boundary: Option<Vec<crate::orderby::OrderKey>> = None;
let order_bound =
crate::orderby::order_by_bound_positions(&order_by, schema_cols, Some(alias));
const BOUNDARY_WINDOW: u32 = 8192;
let mut boundary_checks: u32 = 0;
let mut boundary_rejects: u32 = 0;
let mut boundary_check_on = true;
let mut process_row = |row: &Row<'static>, loop_idx: usize| -> Result<(), EngineError> {
if loop_idx.is_multiple_of(256) {
cancel.check()?;
}
if let Some(cw) = &compiled_where {
let cond = eval::eval_compiled(cw, row, &ctx, &mut eval_stack)
.map_err(EngineError::Eval)?;
if !crate::eval::predicate_is_true(&cond, "WHERE", ctx.mysql_dialect)? {
return Ok(());
}
} else if let Some(where_expr) = &stmt.where_ {
let cond =
self.eval_expr_with_correlated(where_expr, row, &ctx, cancel, Some(&mut memo))?;
if !crate::eval::predicate_is_true(&cond, "WHERE", ctx.mysql_dialect)? {
return Ok(());
}
}
let order_keys = if order_by.is_empty() || stmt.distinct || srf_position.is_some() {
Vec::new()
} else {
let mut buf = key_pool.pop().unwrap_or_default();
crate::orderby::build_order_keys_bound(
&order_by,
&order_bound,
row,
&ctx,
&mut buf,
)?;
if boundary_check_on
&& let Some((_, descs)) = &topk_stream
&& let Some(b) = &topk_boundary
{
boundary_checks += 1;
let loses = crate::orderby::cmp_multi_key_in(&buf, b, descs, &order_colls)
== core::cmp::Ordering::Greater;
if loses {
boundary_rejects += 1;
}
if boundary_checks == BOUNDARY_WINDOW {
boundary_check_on = boundary_rejects.saturating_mul(4) >= boundary_checks;
}
if loses {
buf.clear();
key_pool.push(buf);
return Ok(());
}
}
buf
};
if srf_position.is_some() {
let plan = srf_plan.as_mut().expect("srf_position implies a plan");
for out in expand_srf_row_with(self, plan, &projection, row, &ctx)? {
if stmt.distinct {
let bucket = seen_distinct
.entry(norm_hash_row(&out, &distinct_hb, ctx.mysql_dialect))
.or_default();
if bucket
.iter()
.any(|&i| row_eq_norm(&tagged[i].1, &out, ctx.mysql_dialect))
{
continue;
}
bucket.push(tagged.len());
}
budget.charge(approx_row_bytes(&out))?;
let keys = if order_by.is_empty() {
Vec::new()
} else {
let mut kv: Vec<Value<'static>> = Vec::with_capacity(order_by.len());
for (k, ob) in order_by.iter().enumerate() {
kv.push(match srf_order_cols.get(k).copied().flatten() {
Some(p) => out.values.get(p).cloned().unwrap_or(Value::Null),
None => eval::eval_expr(&ob.expr, row, &ctx)
.map_err(EngineError::Eval)?,
});
}
let key_row = Row::new(kv);
let mut buf = Vec::new();
crate::orderby::build_order_keys_bound(
&order_by,
&srf_key_bound,
&key_row,
&ctx,
&mut buf,
)?;
buf
};
tagged.push((keys, out));
}
} else {
let values = &mut proj_buf;
values.clear();
values.reserve(projection.len());
for (i, p) in projection.iter().enumerate() {
if any_scalarsq_fast && let Some(fp) = &scalarsq_fast[i] {
values.push(self.probe_with_pk_fast_path(fp, row));
continue;
}
if any_proj_const && let Some(v) = &proj_const[i] {
values.push(v.clone());
continue;
}
if any_proj_direct && let Some(pos) = proj_direct[i] {
crate::bump_counter!(crate::select::PROJ_DIRECT_FIRE);
values.push(row.values[pos].clone().into_owned());
continue;
}
let pass_memo = early_cap.is_none_or(|cap| cap > 1000);
let memo_arg = if pass_memo { Some(&mut memo) } else { None };
values.push(
self.eval_expr_with_correlated(&p.expr, row, &ctx, cancel, memo_arg)?,
);
}
crate::bump_counter!(crate::select::PROJ_ROW_BUILT);
if stmt.distinct {
let bucket = seen_distinct
.entry(norm_hash_values(&proj_buf, &distinct_hb, ctx.mysql_dialect))
.or_default();
if bucket
.iter()
.any(|&i| values_eq_norm(&tagged[i].1.values, &proj_buf, ctx.mysql_dialect))
{
crate::bump_counter!(crate::select::DISTINCT_DUP_DROPPED);
return Ok(());
}
bucket.push(tagged.len());
}
let out = Row::new(core::mem::replace(
&mut proj_buf,
proj_pool.pop().unwrap_or_default(),
));
let order_keys = if stmt.distinct && !order_by.is_empty() {
build_order_keys(&order_by, row, &ctx)?
} else {
order_keys
};
budget.charge(approx_row_bytes(&out))?;
tagged.push((order_keys, out));
}
if let Some((k, descs)) = &topk_stream {
crate::orderby::topk_trim_recycling(
&mut tagged,
*k,
descs,
&mut proj_pool,
&mut key_pool,
&mut topk_boundary,
);
}
Ok(())
};
let scan_snapshot = self.current_snapshot();
let mut emitted: usize = 0;
if let Some(rows) = &indexed_rows {
for (loop_idx, cow) in rows.iter().enumerate() {
if let Some(cap) = early_cap
&& emitted >= cap
{
break;
}
process_row(cow.as_ref(), loop_idx)?;
emitted = emitted.saturating_add(1);
}
} else {
let mut rows_cur = table.rows().run_cursor();
for i in 0..table.row_count() {
if let Some(cap) = early_cap
&& emitted >= cap
{
break;
}
if !table.is_row_visible(i, &scan_snapshot) {
continue;
}
let Some(row) = rows_cur.get(i) else { continue };
process_row(row, i)?;
emitted = emitted.saturating_add(1);
}
let cold_rows = self.iter_cold_rows_of_table(table);
for (offset, row) in cold_rows.iter().enumerate() {
if let Some(cap) = early_cap
&& emitted >= cap
{
break;
}
process_row(row, table.row_count() + offset)?;
emitted = emitted.saturating_add(1);
}
}
if !order_by.is_empty() {
let keep = if stmt.limit_with_ties
|| self.env_cfg().disable_topk
{
None
} else {
stmt.limit_literal()
.map(|l| l as usize + stmt.offset_literal().map_or(0, |o| o as usize))
};
let descs: Vec<bool> = order_by.iter().map(|o| o.desc).collect();
crate::orderby::partial_sort_tagged_in(&mut tagged, keep, &descs, &order_colls);
}
let output_rows: Vec<Row<'static>> = if stmt.limit_with_ties && !stmt.distinct {
apply_offset_and_limit_tagged(
&mut tagged,
stmt.offset_literal(),
stmt.limit_literal(),
true,
);
tagged.into_iter().map(|(_, r)| r).collect()
} else {
let mut output_rows: Vec<Row<'static>> = tagged.into_iter().map(|(_, r)| r).collect();
apply_offset_and_limit(
&mut output_rows,
stmt.offset_literal(),
stmt.limit_literal(),
);
output_rows
};
let columns: Vec<ColumnSchema> = projection
.into_iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name, p.ty, p.nullable);
c.user_enum_type = p.user_enum_type;
c.collation_name = p.collation_name;
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
Ok(QueryResult::Rows {
columns,
rows: output_rows,
})
}
fn finish_agg_result(
&self,
mut agg: aggregate::AggResult,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
apply_offset_and_limit(&mut agg.rows, stmt.offset_literal(), stmt.limit_literal());
if !agg.deferred.is_empty() {
apply_offset_and_limit(
&mut agg.synth_rows,
stmt.offset_literal(),
stmt.limit_literal(),
);
let ctx = EvalContext::new(&agg.synth_schema, None);
let mut memo = memoize::MemoizeCache::default();
for (_, expr) in &agg.deferred {
let mut subs: Vec<&SelectStatement> = Vec::new();
collect_scalar_subqueries(expr, &mut subs);
for sub in subs {
let repr = alloc::format!("{sub}");
if memo.group_maps.contains_key(&repr) {
continue;
}
if let Some(gm) = self.try_batch_correlated_scalar(
sub,
Some((&agg.synth_rows, &ctx)),
cancel,
)? {
memo.group_maps.insert(repr, Some(alloc::rc::Rc::new(gm)));
}
}
}
for (ri, srow) in agg.synth_rows.iter().enumerate() {
cancel.check()?;
for (col, expr) in &agg.deferred {
let v =
self.eval_expr_with_correlated(expr, srow, &ctx, cancel, Some(&mut memo))?;
if let Some(cell) = agg.rows[ri].values.get_mut(*col) {
*cell = v;
}
}
}
}
Ok(QueryResult::Rows {
columns: agg.columns,
rows: agg.rows,
})
}
fn try_spill_sorted_scan(
&self,
stmt: &SelectStatement,
from: &FromClause,
cancel: CancelToken<'_>,
) -> Result<Option<QueryResult>, EngineError> {
if !self.can_spill()
|| stmt.order_by.is_empty()
|| stmt.distinct
|| stmt.limit_with_ties
|| stmt.limit_literal().is_some()
|| !from.joins.is_empty()
|| from.primary.lateral_subquery.is_some()
|| from.primary.unnest_expr.is_some()
|| from.primary.generate_series_args.is_some()
|| select_has_window(stmt)
{
return Ok(None);
}
if !from.primary.only
&& crate::partition::has_children(self.active_catalog(), &from.primary.name)
{
return Ok(None);
}
let Some(table) = self.active_catalog().get(&from.primary.name) else {
return Ok(None);
};
if table.has_cold_rows_fast() {
return Ok(None);
}
let alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str());
let cols = table.schema().columns.clone();
let sess = self.dml_session();
let ctx = EvalContext::new(&cols, Some(alias))
.with_catalog(self.active_catalog())
.with_session(&sess);
let projection = build_projection(&stmt.items, &cols, alias, self.backslash_escapes)?;
let order_by = stmt.order_by.clone();
let order_bound = crate::orderby::order_by_bound_positions(&order_by, &cols, Some(alias));
let descs: Vec<bool> = order_by.iter().map(|o| o.desc).collect();
let needed = Self::sort_record_columns_needed(&stmt.items, &order_bound, cols.len(), &ctx);
let mut sorter = crate::extsort::ExternalSorter::new(
self.temp_run_factory,
self.session_work_mem_bytes(),
cols.clone(),
&descs,
)
.with_stats(&self.spill_stats)
.with_pruned(&needed);
let snapshot = self.current_snapshot();
let mut keys: Vec<OrderKey> = Vec::new();
for (i, row) in table.scan_visible_from(0, &snapshot) {
if i.is_multiple_of(256) {
cancel.check()?;
}
if let Some(w) = &stmt.where_ {
let cond = crate::eval::eval_expr(w, row, &ctx).map_err(EngineError::Eval)?;
if !crate::eval::predicate_is_true(&cond, "WHERE", ctx.mysql_dialect)? {
continue;
}
}
keys.clear();
crate::orderby::build_order_keys_bound(&order_by, &order_bound, row, &ctx, &mut keys)?;
sorter.push(&mut keys, row)?;
}
let key_ctx = &ctx;
let rows = sorter.finish(
|src| {
let mut buf = Vec::new();
crate::orderby::build_order_keys_bound(
&order_by,
&order_bound,
src,
key_ctx,
&mut buf,
)?;
Ok(buf)
},
|src| {
let mut values = Vec::with_capacity(projection.len());
for p in &projection {
values.push(
crate::eval::eval_expr(&p.expr, src, key_ctx).map_err(EngineError::Eval)?,
);
}
Ok(Row::new(values))
},
)?;
let columns: Vec<ColumnSchema> = projection
.iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name.clone(), p.ty, p.nullable);
c.user_enum_type = p.user_enum_type.clone();
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
Ok(Some(QueryResult::Rows { columns, rows }))
}
pub(crate) fn sort_record_columns_needed(
items: &[SelectItem],
order_bound: &[Option<usize>],
arity: usize,
ctx: &EvalContext,
) -> Vec<bool> {
let all_bare = items.iter().all(|i| {
matches!(
i,
SelectItem::Expr {
expr: Expr::Column(_),
..
}
)
});
if !all_bare || order_bound.iter().any(Option::is_none) {
return Vec::new();
}
let mut mask = alloc::vec![false; arity];
for item in items {
if let SelectItem::Expr {
expr: Expr::Column(c),
..
} = item
{
match crate::eval::find_column_pos(c, ctx) {
Some(p) if p < arity => mask[p] = true,
_ => return Vec::new(),
}
}
}
for p in order_bound.iter().flatten() {
if *p < arity {
mask[*p] = true;
} else {
return Vec::new();
}
}
mask
}
fn try_spill_sorted_stream<F>(
&self,
stmt: &SelectStatement,
from: &FromClause,
cancel: CancelToken<'_>,
emit: &mut F,
) -> Result<Option<usize>, EngineError>
where
F: FnMut(crate::StreamItem<'_>) -> Result<(), EngineError>,
{
if !self.can_spill()
|| stmt.order_by.is_empty()
|| stmt.distinct
|| stmt.limit_with_ties
|| stmt.limit.is_some()
|| stmt.offset.is_some()
|| stmt.having.is_some()
|| stmt.group_by.is_some()
|| !stmt.unions.is_empty()
|| !from.joins.is_empty()
|| from.primary.lateral_subquery.is_some()
|| from.primary.unnest_expr.is_some()
|| from.primary.as_of_segment.is_some()
|| from.primary.generate_series_args.is_some()
|| select_has_window(stmt)
|| aggregate::uses_aggregate(stmt)
{
return Ok(None);
}
if stmt
.items
.iter()
.any(|i| matches!(i, SelectItem::Expr { expr, .. } if is_top_level_unnest(expr)))
{
return Ok(None);
}
crate::orderby::check_order_by_legality(stmt)?;
crate::orderby::check_order_by_positions(stmt)?;
crate::window::reject_window_in_row_clauses(stmt)?;
if !from.primary.only
&& crate::partition::has_children(self.active_catalog(), &from.primary.name)
{
return Ok(None);
}
let Some(table) = self.active_catalog().get(&from.primary.name) else {
return Ok(None);
};
if table.has_cold_rows_fast() {
return Ok(None);
}
let alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str());
let cols = table.schema().columns.clone();
let sess = self.dml_session();
let ctx = EvalContext::new(&cols, Some(alias))
.with_catalog(self.active_catalog())
.with_session(&sess);
let projection = build_projection(&stmt.items, &cols, alias, self.backslash_escapes)?;
let order_by = stmt.order_by.clone();
let order_bound = crate::orderby::order_by_bound_positions(&order_by, &cols, Some(alias));
let descs: Vec<bool> = order_by.iter().map(|o| o.desc).collect();
let needed = Self::sort_record_columns_needed(&stmt.items, &order_bound, cols.len(), &ctx);
let mut sorter = crate::extsort::ExternalSorter::new(
self.temp_run_factory,
self.session_work_mem_bytes(),
cols.clone(),
&descs,
)
.with_stats(&self.spill_stats)
.with_pruned(&needed);
let snapshot = self.current_snapshot();
let mut keys: Vec<OrderKey> = Vec::new();
for (i, row) in table.scan_visible_from(0, &snapshot) {
if i.is_multiple_of(256) {
cancel.check()?;
}
if let Some(w) = &stmt.where_ {
let cond = crate::eval::eval_expr(w, row, &ctx).map_err(EngineError::Eval)?;
if !crate::eval::predicate_is_true(&cond, "WHERE", ctx.mysql_dialect)? {
continue;
}
}
keys.clear();
crate::orderby::build_order_keys_bound(&order_by, &order_bound, row, &ctx, &mut keys)?;
sorter.push(&mut keys, row)?;
}
let columns: Vec<ColumnSchema> = projection
.iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name.clone(), p.ty, p.nullable);
c.user_enum_type = p.user_enum_type.clone();
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
emit(crate::StreamItem::Header(&columns))?;
let key_ctx = &ctx;
let mut emitted_since_check = 0usize;
let n = sorter.finish_each(
|src| {
let mut buf = Vec::new();
crate::orderby::build_order_keys_bound(
&order_by,
&order_bound,
src,
key_ctx,
&mut buf,
)?;
Ok(buf)
},
|src, values| {
for p in &projection {
values.push(
crate::eval::eval_expr(&p.expr, src, key_ctx).map_err(EngineError::Eval)?,
);
}
Ok(())
},
|cells| {
emitted_since_check += 1;
if emitted_since_check >= 256 {
emitted_since_check = 0;
cancel.check()?;
}
emit(crate::StreamItem::Row(crate::RowCells::Values(cells)))
},
)?;
Ok(Some(n))
}
#[inline]
fn stream_project_row<F>(
row: &spg_storage::Row<'static>,
where_: Option<&Expr>,
projection: &[ProjectedItem],
bound_pos: &[Option<usize>],
ctx: &crate::eval::EvalContext<'_>,
values: &mut Vec<Value<'static>>,
emit: &mut F,
) -> Result<bool, EngineError>
where
F: FnMut(crate::StreamItem<'_>) -> Result<(), EngineError>,
{
if let Some(w) = where_ {
let cond = crate::eval::eval_expr(w, row, ctx).map_err(EngineError::Eval)?;
if !crate::eval::predicate_is_true(&cond, "WHERE", ctx.mysql_dialect)? {
return Ok(false);
}
}
values.clear();
for (p, bound) in projection.iter().zip(bound_pos) {
values.push(match bound {
Some(pos) => crate::eval::column_at(*pos, row, ctx).map_err(EngineError::Eval)?,
None => crate::eval::eval_expr(&p.expr, row, ctx).map_err(EngineError::Eval)?,
});
}
emit(crate::StreamItem::Row(crate::RowCells::Values(values)))?;
Ok(true)
}
fn try_stream_single_table<F>(
&self,
stmt: &SelectStatement,
from: &FromClause,
cancel: CancelToken<'_>,
emit: &mut F,
) -> Result<Option<usize>, EngineError>
where
F: FnMut(crate::StreamItem<'_>) -> Result<(), EngineError>,
{
let Some(table) = self.active_catalog().get(&from.primary.name) else {
return Ok(None);
};
if table.has_cold_rows_fast() {
return Ok(None);
}
let alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str());
let cols = table.schema().columns.clone();
let sess = self.dml_session();
let ctx = EvalContext::new(&cols, Some(alias))
.with_catalog(self.active_catalog())
.with_session(&sess);
let projection = build_projection(&stmt.items, &cols, alias, self.backslash_escapes)?;
let columns: Vec<ColumnSchema> = projection
.iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name.clone(), p.ty, p.nullable);
c.user_enum_type = p.user_enum_type.clone();
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
emit(crate::StreamItem::Header(&columns))?;
let bound_pos: Vec<Option<usize>> = projection
.iter()
.map(|p| match &p.expr {
Expr::Column(c) => match crate::eval::locate_column(c, &ctx) {
Ok(Some(pos)) => Some(pos),
_ => None,
},
_ => None,
})
.collect();
let snapshot = self.current_snapshot();
let seek_positions: Option<Vec<usize>> = stmt.where_.as_ref().and_then(|w| {
crate::index_access::try_index_seek_positions(w, &cols, table, alias, &snapshot)
});
let mut values: Vec<Value<'static>> = Vec::with_capacity(projection.len());
let mut count: usize = 0;
match seek_positions {
Some(mut positions) => {
positions.sort_unstable();
for (n, pos) in positions.into_iter().enumerate() {
if n.is_multiple_of(256) {
cancel.check()?;
}
let Some(row) = table.rows().get(pos) else {
continue;
};
if Self::stream_project_row(
row,
stmt.where_.as_ref(),
&projection,
&bound_pos,
&ctx,
&mut values,
emit,
)? {
count += 1;
}
}
}
None => {
for (i, row) in table.scan_visible_from(0, &snapshot) {
if i.is_multiple_of(256) {
cancel.check()?;
}
if Self::stream_project_row(
row,
stmt.where_.as_ref(),
&projection,
&bound_pos,
&ctx,
&mut values,
emit,
)? {
count += 1;
}
}
}
}
Ok(Some(count))
}
pub(crate) fn try_exec_joined_streaming<F>(
&self,
stmt: &SelectStatement,
cancel: CancelToken<'_>,
emit: &mut F,
) -> Result<Option<usize>, EngineError>
where
F: FnMut(crate::StreamItem<'_>) -> Result<(), EngineError>,
{
let Some(from) = &stmt.from else {
return Ok(None);
};
if self.select_reads_policy_subject_table(stmt) {
return Ok(None);
}
let _single_table = from.joins.is_empty();
if !stmt.order_by.is_empty()
&& from.joins.is_empty()
&& let Some(n) = self.try_spill_sorted_stream(stmt, from, cancel, emit)?
{
return Ok(Some(n));
}
if !stmt.order_by.is_empty()
|| stmt.limit.is_some()
|| stmt.offset.is_some()
|| stmt.having.is_some()
|| stmt.group_by.is_some()
|| stmt.distinct
|| !stmt.unions.is_empty()
|| stmt.limit_with_ties
{
return Ok(None);
}
if aggregate::uses_aggregate(stmt) {
return Ok(None);
}
if select_has_window(stmt) {
return Ok(None);
}
if stmt
.items
.iter()
.any(|i| matches!(i, SelectItem::Expr { expr, .. } if is_top_level_unnest(expr)))
{
return Ok(None);
}
if from.joins.is_empty()
&& from.primary.unnest_expr.is_none()
&& from.primary.lateral_subquery.is_none()
&& from.primary.as_of_segment.is_none()
&& from.primary.generate_series_args.is_none()
&& let Some(n) = self.try_stream_single_table(stmt, from, cancel, emit)?
{
return Ok(Some(n));
}
let mut budget = ByteBudget::new(self.max_query_bytes);
let deferred = {
let mut needed = alloc::collections::BTreeSet::new();
let prunable = collect_qualified_refs(stmt, &mut needed).is_some();
self.build_joined_filtered_rows(
from,
stmt.where_.as_ref(),
cancel,
if prunable { Some(&needed) } else { None },
&mut budget,
)?
};
let combined_schema = &deferred.combined_schema;
let joined_sess = self.dml_session();
let ctx = EvalContext::new(combined_schema, None)
.with_catalog(self.active_catalog())
.with_session(&joined_sess);
let projection =
build_projection(&stmt.items, combined_schema, "", self.backslash_escapes)?;
let bound_pos = |e: &Expr| -> Option<usize> {
match e {
Expr::Column(c) => eval::find_column_pos(c, &ctx),
_ => None,
}
};
let proj_decomposed: Vec<(usize, usize)> = {
let mut out = Vec::with_capacity(projection.len());
for p in &projection {
let Some(abs) = bound_pos(&p.expr) else {
return Ok(None);
};
let Some(k) = deferred
.offsets
.partition_point(|&o| o <= abs)
.checked_sub(1)
else {
return Ok(None);
};
out.push((k, abs - deferred.offsets[k]));
}
out
};
let columns: Vec<ColumnSchema> = projection
.iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name.clone(), p.ty, p.nullable);
c.user_enum_type = p.user_enum_type.clone();
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
emit(crate::StreamItem::Header(&columns))?;
let sources_ref = &deferred.sources;
let stride = deferred.stride;
let survivors_ref = &deferred.survivors;
let n_surv = if stride == 0 {
0
} else {
survivors_ref.len() / stride
};
let null_value = Value::Null;
let mut cell_refs: Vec<&Value> = Vec::with_capacity(projection.len());
let mut count: usize = 0;
for surv_i in 0..n_surv {
if surv_i.is_multiple_of(256) {
cancel.check()?;
}
let tuple = &survivors_ref[surv_i * stride..(surv_i + 1) * stride];
cell_refs.clear();
for &(k, col_in_src) in &proj_decomposed {
let ri = tuple[k];
let v: &Value = if ri == usize::MAX {
&null_value
} else {
sources_ref[k]
.get(ri)
.and_then(|r| r.values.get(col_in_src))
.unwrap_or(&null_value)
};
cell_refs.push(v);
}
emit(crate::StreamItem::Row(crate::RowCells::Refs(&cell_refs)))?;
count += 1;
}
Ok(Some(count))
}
fn exec_joined_select(
&self,
stmt: &SelectStatement,
from: &FromClause,
cancel: CancelToken<'_>,
) -> Result<QueryResult, EngineError> {
if let Some(out) = self.try_count_star_left_anti_join_fast(stmt, from)? {
return Ok(out);
}
if let Some(out) = self.try_streamed_inner_join_walk_topn(stmt, from, cancel)? {
return Ok(out);
}
if let Some(out) = self.try_streamed_inner_join_topn(stmt, from, cancel)? {
return Ok(out);
}
let mut budget = ByteBudget::new(self.max_query_bytes);
let deferred = {
let mut needed = alloc::collections::BTreeSet::new();
let prunable = collect_qualified_refs(stmt, &mut needed).is_some();
self.build_joined_filtered_rows(
from,
stmt.where_.as_ref(),
cancel,
if prunable { Some(&needed) } else { None },
&mut budget,
)?
};
let combined_schema = &deferred.combined_schema;
let joined_sess = self.dml_session();
let ctx = EvalContext::new(combined_schema, None)
.with_catalog(self.active_catalog())
.with_session(&joined_sess);
if aggregate::uses_aggregate(stmt) {
let refs = deferred.row_refs();
let agg_memo = core::cell::RefCell::new(memoize::MemoizeCache::default());
let agg_correlated = |e: &Expr, r: &Row<'static>, c: &EvalContext<'_>| {
self.eval_expr_with_correlated(e, r, c, cancel, Some(&mut agg_memo.borrow_mut()))
.map_err(|err| match err {
EngineError::Eval(ev) => ev,
other => eval::EvalError::TypeMismatch {
detail: alloc::format!("{other}"),
},
})
};
let agg = aggregate::run(
stmt,
crate::join::AggRows::Refs(&refs),
combined_schema,
None,
Some(&agg_correlated),
self.parallel_runner.0.as_deref(),
Some(self.active_catalog()),
Some(self),
)?;
return self.finish_agg_result(agg, stmt, cancel);
}
let projection =
build_projection(&stmt.items, combined_schema, "", self.backslash_escapes)?;
if !self.srf_target_idxs(&projection).is_empty() {
let refs = deferred.row_refs();
let rows: Vec<Row<'static>> = refs.iter().map(|r| r.as_row().into_owned()).collect();
let mut s2 = stmt.clone();
s2.where_ = None;
let schema = combined_schema.clone();
return self.exec_select_over_rows(&s2, rows, schema, "", cancel);
}
let refs = deferred.row_refs();
let bound_pos = |e: &Expr| -> Option<usize> {
match e {
Expr::Column(c) if c.qualifier.is_some() => eval::find_column_pos(c, &ctx),
_ => None,
}
};
let proj_pos: Vec<Option<usize>> = projection.iter().map(|p| bound_pos(&p.expr)).collect();
let all_proj_bound = proj_pos.iter().all(Option::is_some);
let proj_decomposed: Vec<Option<(usize, usize)>> = proj_pos
.iter()
.map(|p| {
p.and_then(|abs| {
let k = deferred
.offsets
.partition_point(|&o| o <= abs)
.checked_sub(1)?;
Some((k, abs - deferred.offsets[k]))
})
})
.collect();
let whole_row_src: Vec<Option<usize>> = projection
.iter()
.map(|p| {
let Expr::Column(c) = &p.expr else {
return None;
};
if !matches!(eval::locate_column(c, &ctx), Ok(None)) {
return None;
}
let prefix = alloc::format!("{name}.", name = c.name);
let abs = deferred
.combined_schema
.iter()
.position(|s| s.name.starts_with(&prefix))?;
deferred
.offsets
.partition_point(|&o| o <= abs)
.checked_sub(1)
})
.collect();
let need_eval_row = !all_proj_bound || !stmt.order_by.is_empty();
let mut tagged: Vec<(Vec<OrderKey>, Row<'static>)> = Vec::new();
let mut proj_memo = memoize::MemoizeCache::default();
let sources_ref = &deferred.sources;
let stride = deferred.stride;
let survivors_ref = &deferred.survivors;
let n_surv = survivors_ref.len() / stride.max(1);
let topk_stream: Option<(usize, Vec<bool>)> = if !stmt.order_by.is_empty()
&& !stmt.distinct
&& !stmt.limit_with_ties
&& !self.env_cfg().disable_topk
{
stmt.limit_literal().and_then(|l| {
let keep = (l as usize).saturating_add(stmt.offset_literal().unwrap_or(0) as usize);
(keep >= 1).then(|| (keep, stmt.order_by.iter().map(|o| o.desc).collect()))
})
} else {
None
};
let mut seen_distinct: hashbrown::HashMap<u64, alloc::vec::Vec<usize>> =
hashbrown::HashMap::new();
let distinct_hb = hashbrown::DefaultHashBuilder::default();
for surv_i in 0..n_surv {
let tuple = &survivors_ref[surv_i * stride..(surv_i + 1) * stride];
let row = &refs[surv_i];
let materialised: Option<Cow<'_, Row<'static>>> = if need_eval_row {
Some(row.as_row())
} else {
None
};
let mut values = Vec::with_capacity(projection.len());
for (i, p) in projection.iter().enumerate() {
if let Some((k, col_in_src)) = proj_decomposed[i] {
let ri = tuple[k];
let v: Value<'static> = if ri == usize::MAX {
Value::Null
} else {
sources_ref[k]
.get(ri)
.and_then(|r| r.values.get(col_in_src))
.cloned()
.map(Value::into_owned)
.unwrap_or(Value::Null)
};
values.push(v);
} else if let Some(pos) = proj_pos[i] {
values.push(
row.get(pos)
.cloned()
.map(Value::into_owned)
.unwrap_or(Value::Null),
);
} else if let Some(k) = whole_row_src[i]
&& tuple[k] == usize::MAX
{
values.push(Value::Null);
} else {
let mrow = materialised.as_deref().expect("materialised for eval");
values.push(self.eval_expr_with_correlated(
&p.expr,
mrow,
&ctx,
cancel,
Some(&mut proj_memo),
)?);
}
}
let out_row = Row::new(values);
if stmt.distinct {
let bucket = seen_distinct
.entry(norm_hash_row(&out_row, &distinct_hb, ctx.mysql_dialect))
.or_default();
if bucket
.iter()
.any(|&i| row_eq_norm(&tagged[i].1, &out_row, ctx.mysql_dialect))
{
continue;
}
bucket.push(tagged.len());
}
let order_keys = if stmt.order_by.is_empty() {
Vec::new()
} else {
let mrow = materialised.as_deref().expect("materialised for order by");
build_order_keys(&stmt.order_by, mrow, &ctx)?
};
budget.charge(approx_row_bytes(&out_row))?;
tagged.push((order_keys, out_row));
if let Some((k, descs)) = &topk_stream {
topk_trim(&mut tagged, *k, descs);
}
}
if !stmt.order_by.is_empty() {
let keep = if self.env_cfg().disable_topk {
None
} else {
stmt.limit_literal()
.map(|l| l as usize + stmt.offset_literal().map_or(0, |o| o as usize))
};
let descs: Vec<bool> = stmt.order_by.iter().map(|o| o.desc).collect();
let colls = crate::orderby::order_by_collations(&stmt.order_by, &ctx)?;
crate::orderby::partial_sort_tagged_in(&mut tagged, keep, &descs, &colls);
}
let mut output_rows: Vec<Row<'static>> = tagged.into_iter().map(|(_, r)| r).collect();
apply_offset_and_limit(
&mut output_rows,
stmt.offset_literal(),
stmt.limit_literal(),
);
let columns: Vec<ColumnSchema> = projection
.into_iter()
.map(|p| {
let mut c = ColumnSchema::new(p.output_name, p.ty, p.nullable);
c.user_enum_type = p.user_enum_type;
c.collation_name = p.collation_name;
c.mysql_fsp = p.mysql_fsp;
c
})
.collect();
Ok(QueryResult::Rows {
columns,
rows: output_rows,
})
}
}
impl Engine {
fn exec_select_as_of_segment(
&self,
stmt: &SelectStatement,
from: &spg_sql::ast::FromClause,
segment_id: u32,
) -> Result<QueryResult, EngineError> {
if !from.joins.is_empty()
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| !stmt.unions.is_empty()
|| !stmt.order_by.is_empty()
|| stmt.offset.is_some()
|| stmt.distinct
|| aggregate::uses_aggregate(stmt)
{
return Err(EngineError::Unsupported(
"AS OF SEGMENT supports SELECT projection + WHERE + LIMIT only \
(joins / aggregates / ORDER BY are STABILITY § \"Out of v6.10\")"
.into(),
));
}
let table = self
.active_catalog()
.get(&from.primary.name)
.ok_or_else(|| StorageError::TableNotFound {
name: from.primary.name.clone(),
})?;
let schema = table.schema().clone();
let schema_cols = &schema.columns;
let alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str());
let ctx = self.ev_ctx(schema_cols, Some(alias));
let seg = self
.active_catalog()
.cold_segment(segment_id)
.ok_or_else(|| {
EngineError::Unsupported(alloc::format!(
"AS OF SEGMENT: cold segment {segment_id} not registered"
))
})?;
let mut out_rows: Vec<Row<'static>> = Vec::new();
let mut limit_remaining: Option<usize> =
stmt.limit_literal().and_then(|n| usize::try_from(n).ok());
for (_key, body) in seg.scan() {
let (row, _consumed) =
spg_storage::decode_row_body_dense(&body, &schema, seg.codec_version())
.map_err(EngineError::Storage)?;
if let Some(where_expr) = &stmt.where_ {
let cond = self.eval_expr_simple(where_expr, &row, &ctx)?;
if !crate::eval::predicate_is_true(&cond, "WHERE", ctx.mysql_dialect)? {
continue;
}
}
let projected = self.project_row_simple(&row, &stmt.items, schema_cols, alias)?;
out_rows.push(projected);
if let Some(rem) = limit_remaining.as_mut() {
if *rem == 0 {
out_rows.pop();
break;
}
*rem -= 1;
}
}
let columns = self.derive_output_columns(&stmt.items, schema_cols, alias);
Ok(QueryResult::Rows {
columns,
rows: out_rows,
})
}
fn eval_expr_simple(
&self,
expr: &Expr,
row: &Row<'static>,
ctx: &EvalContext,
) -> Result<Value<'static>, EngineError> {
let cancel = CancelToken::none();
self.eval_expr_with_correlated(expr, row, ctx, cancel, None)
}
}
#[derive(Debug, Clone)]
pub(crate) struct ProjectedItem {
pub(crate) expr: Expr,
pub(crate) output_name: String,
pub(crate) ty: DataType,
pub(crate) nullable: bool,
pub(crate) user_enum_type: Option<String>,
pub(crate) mysql_fsp: Option<u8>,
pub(crate) collation_name: Option<String>,
}
fn expr_is_aggregate_call(e: &Expr) -> bool {
match e {
Expr::FunctionCall { name, .. } => crate::aggregate::is_aggregate_name(name),
Expr::AggregateOrdered { .. } => true,
_ => false,
}
}
fn collect_agg_exprs(e: &Expr, out: &mut Vec<Expr>) {
if expr_is_aggregate_call(e) {
if !out.iter().any(|x| x == e) {
out.push(e.clone());
}
return;
}
match e {
Expr::Binary { lhs, rhs, .. } => {
collect_agg_exprs(lhs, out);
collect_agg_exprs(rhs, out);
}
Expr::Unary { expr, .. }
| Expr::Cast { expr, .. }
| Expr::IsNull { expr, .. }
| Expr::BoolTest { expr, .. }
| Expr::FieldAccess { base: expr, .. } => collect_agg_exprs(expr, out),
Expr::FunctionCall { args, .. } => {
for a in args {
collect_agg_exprs(a, out);
}
}
Expr::Like { expr, pattern, .. } => {
collect_agg_exprs(expr, out);
collect_agg_exprs(pattern, out);
}
Expr::Extract { source, .. } => collect_agg_exprs(source, out),
Expr::WindowFunction {
args,
partition_by,
order_by,
..
} => {
for a in args {
collect_agg_exprs(a, out);
}
for p in partition_by {
collect_agg_exprs(p, out);
}
for (o, _, _) in order_by {
collect_agg_exprs(o, out);
}
}
_ => {}
}
}
fn replace_agg_exprs(e: &mut Expr, aggs: &[Expr]) {
if expr_is_aggregate_call(e) {
if let Some(idx) = aggs.iter().position(|x| x == e) {
*e = Expr::Column(ColumnName {
qualifier: None,
name: alloc::format!("__agg{idx}"),
});
}
return;
}
match e {
Expr::Binary { lhs, rhs, .. } => {
replace_agg_exprs(lhs, aggs);
replace_agg_exprs(rhs, aggs);
}
Expr::Unary { expr, .. }
| Expr::Cast { expr, .. }
| Expr::IsNull { expr, .. }
| Expr::BoolTest { expr, .. }
| Expr::FieldAccess { base: expr, .. } => replace_agg_exprs(expr, aggs),
Expr::FunctionCall { args, .. } => {
for a in args {
replace_agg_exprs(a, aggs);
}
}
Expr::Like { expr, pattern, .. } => {
replace_agg_exprs(expr, aggs);
replace_agg_exprs(pattern, aggs);
}
Expr::Extract { source, .. } => replace_agg_exprs(source, aggs),
Expr::WindowFunction {
args,
partition_by,
order_by,
..
} => {
for a in args {
replace_agg_exprs(a, aggs);
}
for p in partition_by {
replace_agg_exprs(p, aggs);
}
for (o, _, _) in order_by {
replace_agg_exprs(o, aggs);
}
}
_ => {}
}
}
fn rewrite_agg_before_window(stmt: &SelectStatement) -> Option<SelectStatement> {
if !(crate::aggregate::uses_aggregate(stmt) || stmt.group_by.is_some()) {
return None;
}
if !stmt.unions.is_empty() {
return None;
}
let group_cols: Vec<Expr> = stmt.group_by.clone().unwrap_or_default();
if group_cols.iter().any(|g| !matches!(g, Expr::Column(_))) {
return None;
}
stmt.from.as_ref()?;
let mut aggs: Vec<Expr> = Vec::new();
for item in &stmt.items {
if let SelectItem::Expr { expr, .. } = item {
collect_agg_exprs(expr, &mut aggs);
}
}
for ob in &stmt.order_by {
collect_agg_exprs(&ob.expr, &mut aggs);
}
let mut inner_items: Vec<SelectItem> = Vec::new();
for g in &group_cols {
inner_items.push(SelectItem::Expr {
expr: g.clone(),
alias: None,
});
}
for (i, a) in aggs.iter().enumerate() {
inner_items.push(SelectItem::Expr {
expr: a.clone(),
alias: Some(alloc::format!("__agg{i}")),
});
}
let inner = SelectStatement {
items: inner_items,
distinct: false,
distinct_on: Vec::new(),
unions: Vec::new(),
order_by: Vec::new(),
limit: None,
offset: None,
limit_with_ties: false,
window_check_exprs: Vec::new(),
..stmt.clone()
};
let derived = TableRef {
name: "__aggwin".into(),
alias: Some("__aggwin".into()),
only: false,
as_of_segment: None,
unnest_expr: None,
unnest_column_aliases: Vec::new(),
with_ordinality: false,
generate_series_args: None,
lateral_subquery: Some(alloc::boxed::Box::new(inner)),
jsonb_each_text_arg: None,
table_fn_call: None,
rows_from: None,
json_table: None,
scalar_fn_item: false,
};
let mut outer_items = stmt.items.clone();
for item in &mut outer_items {
if let SelectItem::Expr { expr, alias } = item {
if alias.is_none()
&& let Expr::FunctionCall { name, .. } = expr
&& crate::aggregate::is_aggregate_name(name)
{
*alias = Some(name.to_ascii_lowercase());
}
replace_agg_exprs(expr, &aggs);
}
}
let mut outer_order = stmt.order_by.clone();
for ob in &mut outer_order {
replace_agg_exprs(&mut ob.expr, &aggs);
}
let mut outer_distinct_on = stmt.distinct_on.clone();
for e in &mut outer_distinct_on {
replace_agg_exprs(e, &aggs);
}
Some(SelectStatement {
locking: None,
ctes: Vec::new(),
distinct: stmt.distinct,
distinct_on: outer_distinct_on,
items: outer_items,
from: Some(FromClause {
primary: derived,
joins: Vec::new(),
}),
where_: None,
group_by: None,
group_by_all: false,
having: None,
unions: Vec::new(),
order_by: outer_order,
limit: stmt.limit.clone(),
offset: stmt.offset.clone(),
limit_with_ties: stmt.limit_with_ties,
window_check_exprs: Vec::new(),
})
}
struct PeerIndex<'r> {
bh: hashbrown::DefaultHashBuilder,
buckets: hashbrown::HashMap<u64, Vec<usize>>,
rows: &'r [Row<'static>],
mysql: bool,
}
impl<'r> PeerIndex<'r> {
fn build(rows: &'r [Row<'static>], mysql: bool) -> Self {
let bh = hashbrown::DefaultHashBuilder::default();
let mut buckets: hashbrown::HashMap<u64, Vec<usize>> =
hashbrown::HashMap::with_capacity(rows.len());
for (i, r) in rows.iter().enumerate() {
buckets
.entry(norm_hash_row(r, &bh, mysql))
.or_default()
.push(i);
}
Self {
bh,
buckets,
rows,
mysql,
}
}
fn contains(&self, r: &Row<'static>) -> bool {
let h = norm_hash_row(r, &self.bh, self.mysql);
self.buckets
.get(&h)
.is_some_and(|b| b.iter().any(|&i| row_eq_norm(&self.rows[i], r, self.mysql)))
}
fn take_one(&mut self, r: &Row<'static>) -> bool {
let h = norm_hash_row(r, &self.bh, self.mysql);
let Some(b) = self.buckets.get_mut(&h) else {
return false;
};
let Some(pos) = b
.iter()
.position(|&i| row_eq_norm(&self.rows[i], r, self.mysql))
else {
return false;
};
b.swap_remove(pos);
true
}
}
pub(crate) fn dedup_rows(rows: Vec<Row<'static>>, mysql: bool) -> Vec<Row<'static>> {
dedup_by_row(rows, |r| r, mysql)
}
fn dedup_by_row<T>(items: Vec<T>, row_of: impl Fn(&T) -> &Row<'static>, mysql: bool) -> Vec<T> {
if items.len() <= 32 {
let mut out: Vec<T> = Vec::with_capacity(items.len());
for it in items {
if !out
.iter()
.any(|seen| row_eq_norm(row_of(seen), row_of(&it), mysql))
{
out.push(it);
}
}
return out;
}
let bh = hashbrown::DefaultHashBuilder::default();
let mut out: Vec<T> = Vec::with_capacity(items.len().min(1024));
let mut buckets: hashbrown::HashMap<u64, alloc::vec::Vec<usize>> =
hashbrown::HashMap::with_capacity(items.len());
for it in items {
let h = norm_hash_row(row_of(&it), &bh, mysql);
let bucket = buckets.entry(h).or_default();
if !bucket
.iter()
.any(|&i| row_eq_norm(row_of(&out[i]), row_of(&it), mysql))
{
bucket.push(out.len());
out.push(it);
}
}
out
}
fn norm_hash_row(row: &Row<'static>, bh: &hashbrown::DefaultHashBuilder, mysql: bool) -> u64 {
norm_hash_values(&row.values, bh, mysql)
}
fn norm_hash_values(
values: &[Value<'static>],
bh: &hashbrown::DefaultHashBuilder,
mysql: bool,
) -> u64 {
use core::hash::{BuildHasher, Hash, Hasher};
let mut h = bh.build_hasher();
for v in values {
if mysql {
if let Some(folded) = mysql_dedup_fold(v) {
folded.hash(&mut h);
continue;
}
}
norm_hash_value(v, &mut h);
}
h.finish()
}
fn norm_hash_value<H: core::hash::Hasher>(v: &Value<'static>, h: &mut H) {
const TAG_NULL: u8 = 0;
const TAG_BOOL: u8 = 1;
const TAG_NUM_I64: u8 = 2;
const TAG_NUM_F64: u8 = 3;
const TAG_TEXT: u8 = 4;
const TAG_DATE: u8 = 6;
const TAG_TIME: u8 = 7;
const TAG_TIMESTAMP: u8 = 8;
const TAG_TIMETZ: u8 = 10;
const TAG_UUID: u8 = 11;
const TAG_MONEY: u8 = 12;
const TAG_BYTES: u8 = 13;
const TAG_INTERVAL: u8 = 14;
const TAG_CHAR1: u8 = 15;
const TAG_OPAQUE: u8 = 255;
let num_f64 = |h: &mut H, x: f64| {
if x.is_nan() {
h.write_u8(TAG_NUM_F64);
h.write_u64(0x7ff8_dead_beef_0001); return;
}
const TWO63: f64 = 9_223_372_036_854_775_808.0;
if (-TWO63..TWO63).contains(&x) {
#[allow(clippy::cast_possible_truncation)]
let n = x as i64;
#[allow(clippy::cast_precision_loss)]
if (n as f64) == x {
h.write_u8(TAG_NUM_I64);
h.write_i64(n);
return;
}
}
h.write_u8(TAG_NUM_F64);
h.write_u64(x.to_bits());
};
match v {
Value::Null => h.write_u8(TAG_NULL),
Value::Bool(b) => {
h.write_u8(TAG_BOOL);
h.write_u8(u8::from(*b));
}
Value::SmallInt(n) => {
h.write_u8(TAG_NUM_I64);
h.write_i64(i64::from(*n));
}
Value::Int(n) => {
h.write_u8(TAG_NUM_I64);
h.write_i64(i64::from(*n));
}
Value::BigInt(n) => {
h.write_u8(TAG_NUM_I64);
h.write_i64(*n);
}
Value::Float(x) => num_f64(h, *x),
Value::Numeric {
scaled,
scale,
kind,
} => match kind {
spg_storage::NumericKind::NaN => num_f64(h, f64::NAN),
spg_storage::NumericKind::PosInf => num_f64(h, f64::INFINITY),
spg_storage::NumericKind::NegInf => num_f64(h, f64::NEG_INFINITY),
spg_storage::NumericKind::Finite => {
let (mut s, mut sc) = (*scaled, *scale);
while sc > 0 && s % 10 == 0 {
s /= 10;
sc -= 1;
}
if sc == 0 {
if let Ok(n) = i64::try_from(s) {
h.write_u8(TAG_NUM_I64);
h.write_i64(n);
} else {
num_f64(h, crate::orderby::numeric_to_f64(s, 0));
}
} else {
num_f64(h, crate::orderby::numeric_to_f64(s, sc));
}
}
},
Value::NumericBig(b) => match b.to_i128() {
Some(s) => norm_hash_value(
&Value::Numeric {
scaled: s,
scale: b.scale(),
kind: spg_storage::NumericKind::Finite,
},
h,
),
None => h.write_u8(TAG_OPAQUE),
},
Value::Text(s) | Value::BpChar(s) => {
h.write_u8(TAG_TEXT);
h.write(s.trim_end_matches(' ').as_bytes());
}
Value::Char1(c) => {
h.write_u8(TAG_CHAR1);
h.write_u8(*c);
}
Value::Date(d) => {
h.write_u8(TAG_DATE);
h.write_i32(*d);
}
Value::Time(t) => {
h.write_u8(TAG_TIME);
h.write_i64(*t);
}
Value::Timestamp(t) => {
h.write_u8(TAG_TIMESTAMP);
h.write_i64(*t);
}
Value::TimeTz { us, offset_secs } => {
h.write_u8(TAG_TIMETZ);
h.write_i64(*us);
h.write_i32(*offset_secs);
}
Value::Uuid(u) => {
h.write_u8(TAG_UUID);
h.write(u);
}
Value::Money(c) => {
h.write_u8(TAG_MONEY);
h.write_i64(*c);
}
Value::Bytes(b) => {
h.write_u8(TAG_BYTES);
h.write(b.as_ref());
}
Value::Interval {
months,
days,
micros,
} => {
h.write_u8(TAG_INTERVAL);
h.write_i32(*months);
h.write_i32(*days);
h.write_i64(*micros);
}
Value::Real(x) => num_f64(h, f64::from(*x)),
_ => h.write_u8(TAG_OPAQUE),
}
}
fn mysql_dedup_fold(v: &Value) -> Option<String> {
match v {
Value::Text(s) | Value::BpChar(s) => {
Some(spg_storage::mysql_ci_fold(s.trim_end_matches(' ')))
}
_ => None,
}
}
pub static SCAN_PATH_ENTERED: core::sync::atomic::AtomicU64 = core::sync::atomic::AtomicU64::new(0);
pub static PROJ_DIRECT_FIRE: core::sync::atomic::AtomicU64 = core::sync::atomic::AtomicU64::new(0);
pub static PROJ_ROW_BUILT: core::sync::atomic::AtomicU64 = core::sync::atomic::AtomicU64::new(0);
pub static DISTINCT_DUP_DROPPED: core::sync::atomic::AtomicU64 =
core::sync::atomic::AtomicU64::new(0);
pub(crate) fn row_eq_norm(a: &Row<'static>, b: &Row<'static>, mysql: bool) -> bool {
values_eq_norm(&a.values, &b.values, mysql)
}
pub(crate) fn values_eq_norm(a: &[Value<'static>], b: &[Value<'static>], mysql: bool) -> bool {
a.len() == b.len()
&& a.iter().zip(b).all(|(x, y)| {
if mysql {
if let (Some(fx), Some(fy)) = (mysql_dedup_fold(x), mysql_dedup_fold(y)) {
return fx == fy;
}
}
crate::orderby::value_cmp(x, y) == core::cmp::Ordering::Equal
})
}
pub(crate) fn value_to_order_key(v: &Value) -> Result<OrderKey, EngineError> {
if let Value::Text(s) = v {
return Ok(OrderKey::Text(s.as_ref().into()));
}
if let Value::BpChar(s) = v {
return Ok(OrderKey::Text(s.trim_end_matches(' ').into()));
}
if let Value::Json(s) = v {
return Ok(match crate::json::parse(s) {
Ok(jv) => OrderKey::Json(jv),
Err(_) => OrderKey::Text(s.as_ref().into()),
});
}
match v {
Value::Bytes(b) => return Ok(OrderKey::Bytes(b.as_ref().to_vec())),
Value::NumericBig(b) => return Ok(OrderKey::BigNum((**b).clone())),
Value::Uuid(u) => return Ok(OrderKey::Bytes(u.to_vec())),
Value::Macaddr(m) => return Ok(OrderKey::Bytes(m.to_vec())),
Value::Macaddr8(m) => return Ok(OrderKey::Bytes(m.to_vec())),
Value::PgLsn(l) => return Ok(OrderKey::Bytes(l.to_be_bytes().to_vec())),
Value::Inet { family, bits, addr } | Value::Cidr { family, bits, addr } => {
let mut key = alloc::vec::Vec::with_capacity(18);
key.push(*family);
key.extend_from_slice(addr);
key.push(*bits);
return Ok(OrderKey::Bytes(key));
}
_ => {}
}
let inf = || OrderKey::NullBig;
let arr = match v {
Value::IntArray(a) => Some(
a.iter()
.map(|o| o.map_or_else(inf, |n| OrderKey::Int(i128::from(n))))
.collect(),
),
Value::SmallIntArray(a) => Some(
a.iter()
.map(|o| o.map_or_else(inf, |n| OrderKey::Int(i128::from(n))))
.collect(),
),
Value::BigIntArray(a) => Some(
a.iter()
.map(|o| o.map_or_else(inf, |n| OrderKey::Int(i128::from(n))))
.collect(),
),
Value::BoolArray(a) => Some(
a.iter()
.map(|o| o.map_or_else(inf, |b| OrderKey::Int(i128::from(b))))
.collect(),
),
Value::TextArray(a) => Some(
a.iter()
.map(|o| o.as_ref().map_or_else(inf, |s| OrderKey::Text(s.clone())))
.collect(),
),
#[allow(clippy::cast_precision_loss)]
Value::FloatArray(a) => Some(
a.iter()
.map(|o| o.map_or(OrderKey::NullBig, OrderKey::Num))
.collect(),
),
Value::NumericArray(a) => Some(
a.iter()
.map(|o| {
o.map_or_else(inf, |(m, s)| {
OrderKey::Num(crate::orderby::numeric_to_f64(m, s))
})
})
.collect(),
),
Value::DateArray(a) => Some(
a.iter()
.map(|o| o.map_or_else(inf, |n| OrderKey::Int(i128::from(n))))
.collect(),
),
_ => None,
};
if let Some(elements) = arr {
return Ok(OrderKey::Array(elements));
}
if let Value::Composite(fields) = v {
let elements = fields
.iter()
.map(|(_, fv)| value_to_order_key(fv))
.collect::<Result<alloc::vec::Vec<_>, _>>()?;
return Ok(OrderKey::Array(elements));
}
match v {
Value::SmallInt(n) => return Ok(OrderKey::Int(i128::from(*n))),
Value::Int(n) => return Ok(OrderKey::Int(i128::from(*n))),
Value::BigInt(n) => return Ok(OrderKey::Int(i128::from(*n))),
Value::Date(d) => return Ok(OrderKey::Int(i128::from(*d))),
Value::Timestamp(t) => return Ok(OrderKey::Int(i128::from(*t))),
Value::Time(us) => return Ok(OrderKey::Int(i128::from(*us))),
Value::Year(y) => return Ok(OrderKey::Int(i128::from(*y))),
Value::TimeTz { us, offset_secs } => {
return Ok(OrderKey::Int(
i128::from(*us) - i128::from(*offset_secs) * 1_000_000,
));
}
Value::Money(c) => return Ok(OrderKey::Int(i128::from(*c))),
_ => {}
}
let num = match v {
Value::Null => return Ok(OrderKey::NullBig),
Value::Range { .. } => {
return Err(EngineError::Unsupported(
"ORDER BY of a range value is not supported in v7.17.0".into(),
));
}
Value::Hstore(_) => {
return Err(EngineError::Unsupported(
"ORDER BY of a hstore value is not supported".into(),
));
}
Value::IntArray2D(_) | Value::BigIntArray2D(_) | Value::TextArray2D(_) => {
return Err(EngineError::Unsupported(
"ORDER BY of a 2D array is not supported in v7.17.0".into(),
));
}
#[allow(clippy::cast_precision_loss)]
Value::Numeric { scaled, scale, .. } => {
let mut divisor = 1.0_f64;
for _ in 0..*scale {
divisor *= 10.0;
}
(*scaled as f64) / divisor
}
Value::Float(x) => *x,
Value::Real(x) => f64::from(*x),
Value::Bool(b) => {
if *b {
1.0
} else {
0.0
}
}
Value::Vector(_) | Value::Sq8Vector(_) | Value::HalfVector(_) => {
return Err(EngineError::Unsupported(
"ORDER BY of a raw vector column is not meaningful — use `<->`".into(),
));
}
#[allow(clippy::cast_precision_loss)]
Value::Interval {
months,
days,
micros,
} => {
let total = i128::from(*months) * 30 * 86_400_000_000
+ i128::from(*days) * 86_400_000_000
+ i128::from(*micros);
total as f64
}
Value::Json(_) => {
return Err(EngineError::Unsupported(
"ORDER BY of a JSON value is not supported — cast the document to text first"
.into(),
));
}
_ => {
return Err(EngineError::Unsupported(
"ORDER BY of this value type is not supported".into(),
));
}
};
Ok(OrderKey::Num(num))
}
pub(crate) const CTID_COLUMN: &str = "ctid";
pub(crate) const SYSTEM_COLUMNS: [&str; 6] = ["ctid", "xmin", "xmax", "cmin", "cmax", "tableoid"];
pub(crate) fn is_system_column(name: &str) -> bool {
SYSTEM_COLUMNS.iter().any(|s| name.eq_ignore_ascii_case(s))
}
fn system_column_tail_start(cols: &[ColumnSchema]) -> Option<usize> {
let start = cols.len().checked_sub(SYSTEM_COLUMNS.len())?;
cols[start..]
.iter()
.zip(SYSTEM_COLUMNS)
.all(|(c, name)| c.name.eq_ignore_ascii_case(name))
.then_some(start)
}
fn synthetic_system_positions(cols: &[ColumnSchema]) -> alloc::vec::Vec<bool> {
let mut skip = alloc::vec![false; cols.len()];
fn qualifier(n: &str) -> Option<&str> {
n.rsplit_once('.').map(|(q, _)| q)
}
fn bare(n: &str) -> &str {
n.rsplit('.').next().unwrap_or(n)
}
let mut i = 0;
while i < cols.len() {
let q = qualifier(&cols[i].name);
let mut end = i;
while end < cols.len() && qualifier(&cols[end].name) == q {
end += 1;
}
if let Some(start) = (end - i)
.checked_sub(SYSTEM_COLUMNS.len())
.map(|off| i + off)
&& cols[start..end]
.iter()
.zip(SYSTEM_COLUMNS)
.all(|(c, name)| bare(&c.name).eq_ignore_ascii_case(name))
{
for s in skip.iter_mut().take(end).skip(start) {
*s = true;
}
}
i = end;
}
skip
}
pub(crate) fn expr_references_ctid(e: &Expr) -> bool {
let mut found = false;
crate::expr_analysis::visit_expr_columns_and_subqueries(
e,
&mut |c| {
if is_system_column(&c.name) {
found = true;
}
},
&mut |_| {},
);
found
}
fn references_ctid(stmt: &SelectStatement) -> bool {
let in_expr = expr_references_ctid;
stmt.items.iter().any(|i| match i {
SelectItem::Expr { expr, .. } => in_expr(expr),
_ => false,
}) || stmt.where_.as_ref().is_some_and(in_expr)
|| stmt.order_by.iter().any(|o| in_expr(&o.expr))
|| stmt
.group_by
.as_ref()
.is_some_and(|g| g.iter().any(in_expr))
|| stmt.having.as_ref().is_some_and(in_expr)
}
fn whole_row_projection_schema(alias: &str) -> ColumnSchema {
let mut s = ColumnSchema::new(
alloc::string::String::from(alias),
spg_storage::DataType::Jsonb,
true,
);
s.user_composite_type = Some(alloc::string::String::from(alias));
s
}
pub(crate) fn resolve_projection_column<'a>(
c: &ColumnName,
schema_cols: &'a [ColumnSchema],
table_alias: &str,
) -> Result<Cow<'a, ColumnSchema>, EngineError> {
if let Some(q) = &c.qualifier {
let composite = alloc::format!("{q}.{name}", name = c.name);
if let Some(s) = schema_cols.iter().find(|s| s.name == composite) {
return Ok(Cow::Borrowed(s));
}
if q == table_alias
&& let Some(s) = schema_cols.iter().find(|s| s.name == c.name)
{
return Ok(Cow::Borrowed(s));
}
let prefix = alloc::format!("{q}.");
let qualifier_known =
q == table_alias || schema_cols.iter().any(|s| s.name.starts_with(&prefix));
if !qualifier_known {
return Err(EngineError::Eval(EvalError::UnknownQualifier {
qualifier: q.clone(),
}));
}
return Err(EngineError::Eval(EvalError::ColumnNotFound {
name: c.name.clone(),
}));
}
if let Some(s) = schema_cols.iter().find(|s| s.name == c.name) {
return Ok(Cow::Borrowed(s));
}
let suffix = alloc::format!(".{name}", name = c.name);
let mut matches = schema_cols.iter().filter(|s| s.name.ends_with(&suffix));
let first = matches.next();
let extra = matches.next();
match (first, extra) {
(Some(s), None) => Ok(Cow::Borrowed(s)),
(Some(_), Some(_)) => Err(EngineError::Eval(EvalError::TypeMismatch {
detail: alloc::format!("column reference \"{}\" is ambiguous", c.name),
})),
_ if !table_alias.is_empty() && c.name == table_alias => {
Ok(Cow::Owned(whole_row_projection_schema(table_alias)))
}
_ if table_alias.is_empty() && {
let prefix = alloc::format!("{name}.", name = c.name);
schema_cols.iter().any(|s| s.name.starts_with(&prefix))
} =>
{
Ok(Cow::Owned(whole_row_projection_schema(&c.name)))
}
_ => Err(EngineError::Eval(EvalError::ColumnNotFound {
name: c.name.clone(),
})),
}
}
fn apply_deferred_limit(
rows: alloc::vec::Vec<Row<'static>>,
deferred: &(
Option<spg_sql::ast::LimitExpr>,
Option<spg_sql::ast::LimitExpr>,
),
) -> alloc::vec::Vec<Row<'static>> {
let count = |e: &Option<spg_sql::ast::LimitExpr>| match e {
Some(spg_sql::ast::LimitExpr::Literal(n)) => Some(*n as usize),
_ => None,
};
let mut rows = rows;
if let Some(off) = count(&deferred.1) {
rows = rows.split_off(off.min(rows.len()));
}
if let Some(lim) = count(&deferred.0) {
rows.truncate(lim);
}
rows
}
fn strip_synthetic_order_cols(result: QueryResult) -> QueryResult {
let QueryResult::Rows { columns, rows } = result else {
return result;
};
if !columns.iter().any(|c| c.name.starts_with("__grp_ord_")) {
return QueryResult::Rows { columns, rows };
}
let keep: Vec<usize> = columns
.iter()
.enumerate()
.filter(|(_, c)| !c.name.starts_with("__grp_ord_"))
.map(|(i, _)| i)
.collect();
let new_cols: Vec<ColumnSchema> = keep.iter().map(|&i| columns[i].clone()).collect();
let new_rows: Vec<Row<'static>> = rows
.into_iter()
.map(|r| Row::new(keep.iter().map(|&i| r.values[i].clone()).collect()))
.collect();
QueryResult::Rows {
columns: new_cols,
rows: new_rows,
}
}
#[inline(never)]
fn bind_direct_columns(
projection: &[ProjectedItem],
ctx: &eval::EvalContext<'_>,
) -> Vec<Option<usize>> {
projection
.iter()
.map(|p| match &p.expr {
Expr::Column(c) => eval::compile_column_pos(c, ctx).filter(|pos| {
ctx.columns
.get(*pos)
.is_none_or(|sc| sc.user_composite_type.is_none())
}),
_ => None,
})
.collect()
}
pub(crate) fn default_output_name(expr: &Expr, mysql: bool) -> String {
if mysql {
return expr.to_string();
}
spg_sql::ast::figure_column_name(expr).unwrap_or_else(|| "?column?".to_string())
}
pub(crate) fn build_projection(
items: &[SelectItem],
schema_cols: &[ColumnSchema],
table_alias: &str,
mysql: bool,
) -> Result<Vec<ProjectedItem>, EngineError> {
build_projection_hiding_tail(items, schema_cols, table_alias, mysql, 0)
}
pub(crate) fn build_projection_hiding_tail(
items: &[SelectItem],
schema_cols: &[ColumnSchema],
table_alias: &str,
mysql: bool,
hidden_tail: usize,
) -> Result<Vec<ProjectedItem>, EngineError> {
let visible = schema_cols.len().saturating_sub(hidden_tail);
let joined_schema = table_alias.is_empty()
&& !schema_cols.is_empty()
&& schema_cols.iter().all(|c| c.name.contains('.'));
let bare_name = |name: &str| -> String {
if !joined_schema {
return name.to_string();
}
match name.split_once('.') {
Some((_, rest)) if !rest.is_empty() => rest.to_string(),
_ => name.to_string(),
}
};
let mut out = Vec::new();
for item in items {
match item {
SelectItem::Wildcard => {
let sys_skip = synthetic_system_positions(schema_cols);
for (idx, col) in schema_cols.iter().enumerate() {
if sys_skip[idx] || idx >= visible {
continue;
}
out.push(ProjectedItem {
expr: Expr::Column(ColumnName {
qualifier: None,
name: col.name.clone(),
}),
output_name: bare_name(&col.name),
ty: col.ty,
nullable: col.nullable,
user_enum_type: col.user_enum_type.clone(),
mysql_fsp: col.mysql_fsp,
collation_name: col.collation_name.clone(),
});
}
}
SelectItem::QualifiedWildcard(q) => {
let prefix = alloc::format!("{q}.");
let single_table = !table_alias.is_empty() && q == table_alias;
let mut matched = 0usize;
for col in &schema_cols[..visible] {
let belongs =
col.name.starts_with(&prefix) || (single_table && !col.name.contains('.'));
if !belongs {
continue;
}
matched += 1;
let output_name = col
.name
.strip_prefix(&prefix)
.unwrap_or(&col.name)
.to_string();
out.push(ProjectedItem {
expr: Expr::Column(ColumnName {
qualifier: None,
name: col.name.clone(),
}),
output_name,
ty: col.ty,
nullable: col.nullable,
user_enum_type: col.user_enum_type.clone(),
mysql_fsp: col.mysql_fsp,
collation_name: col.collation_name.clone(),
});
}
if matched == 0 {
return Err(EngineError::Eval(EvalError::UnknownQualifier {
qualifier: q.clone(),
}));
}
}
SelectItem::Expr { expr, alias } => {
if let Expr::Column(c) = expr {
let sch = resolve_projection_column(c, schema_cols, table_alias)?;
let output_name = alias.clone().unwrap_or_else(|| c.name.clone());
out.push(ProjectedItem {
expr: expr.clone(),
output_name,
ty: sch.ty,
nullable: sch.nullable,
user_enum_type: sch.user_enum_type.clone(),
mysql_fsp: sch.mysql_fsp,
collation_name: sch.collation_name.clone(),
});
} else if let Some(shape) = describe::describe_expr(expr, schema_cols) {
let output_name = alias
.clone()
.unwrap_or_else(|| default_output_name(expr, mysql));
out.push(ProjectedItem {
expr: expr.clone(),
output_name,
ty: shape.ty,
nullable: shape.nullable,
user_enum_type: None,
mysql_fsp: crate::eval::expr_mysql_fsp(expr, schema_cols),
collation_name: match expr {
Expr::Column(c) => schema_cols
.iter()
.find(|sc| sc.name.eq_ignore_ascii_case(&c.name))
.and_then(|sc| sc.collation_name.clone()),
_ => None,
},
});
} else {
let output_name = alias
.clone()
.unwrap_or_else(|| default_output_name(expr, mysql));
out.push(ProjectedItem {
expr: expr.clone(),
output_name,
ty: DataType::Text,
nullable: true,
user_enum_type: crate::eval::expr_enum_type_name_pub(expr, schema_cols)
.map(alloc::string::String::from),
mysql_fsp: crate::eval::expr_mysql_fsp(expr, schema_cols),
collation_name: match expr {
Expr::Column(c) => schema_cols
.iter()
.find(|sc| sc.name.eq_ignore_ascii_case(&c.name))
.and_then(|sc| sc.collation_name.clone()),
_ => None,
},
});
}
}
}
}
Ok(out)
}
pub(crate) fn infer_column_types(
columns: &[ColumnSchema],
rows: &[Row<'static>],
) -> Vec<ColumnSchema> {
let mut out = columns.to_vec();
for (col_idx, col) in out.iter_mut().enumerate() {
if col.ty != DataType::Text {
continue;
}
let mut inferred: Option<DataType> = None;
let mut all_null = true;
for row in rows {
let Some(v) = row.values.get(col_idx) else {
continue;
};
let ty = match v {
Value::Null => continue,
Value::SmallInt(_) => DataType::SmallInt,
Value::Int(_) => DataType::Int,
Value::BigInt(_) => DataType::BigInt,
Value::Float(_) => DataType::Float,
Value::Bool(_) => DataType::Bool,
Value::Vector(_) => DataType::Vector {
dim: 0,
encoding: VecEncoding::F32,
},
Value::TextArray(_) => DataType::TextArray,
Value::IntArray(_) => DataType::IntArray,
Value::BigIntArray(_) => DataType::BigIntArray,
Value::SmallIntArray(_) => DataType::SmallIntArray,
Value::FloatArray(_) => DataType::FloatArray,
Value::BoolArray(_) => DataType::BoolArray,
Value::Interval { .. } => DataType::Interval,
_ => DataType::Text,
};
all_null = false;
inferred = Some(match inferred {
None => ty,
Some(prev) if prev == ty => prev,
Some(_) => DataType::Text,
});
}
if let Some(t) = inferred {
col.ty = t;
col.nullable = true;
} else if all_null {
col.nullable = true;
}
}
out
}
fn numeric_rank(t: DataType) -> Option<u8> {
match t {
DataType::SmallInt => Some(1),
DataType::Int => Some(2),
DataType::BigInt => Some(3),
DataType::Numeric { .. } => Some(4),
DataType::Float => Some(5),
_ => None,
}
}
fn resolve_union_common_type(types: &[DataType]) -> Option<DataType> {
if types.len() < 2 {
return None;
}
if types.iter().all(|t| numeric_rank(*t).is_some()) {
return types
.iter()
.max_by_key(|t| numeric_rank(**t).unwrap_or(0))
.copied();
}
let non_text: Vec<&DataType> = types
.iter()
.filter(|t| !matches!(t, DataType::Text))
.collect();
if non_text.iter().all(|t| {
matches!(
t,
DataType::Date | DataType::Timestamp | DataType::Timestamptz
)
}) && non_text
.iter()
.any(|t| matches!(t, DataType::Timestamp | DataType::Timestamptz))
{
if non_text.iter().any(|t| matches!(t, DataType::Timestamptz)) {
return Some(DataType::Timestamptz);
}
return Some(DataType::Timestamp);
}
if non_text.len() == 1 {
return Some(*non_text[0]);
}
if !non_text.is_empty() && non_text.len() < types.len() {
let concrete: Vec<DataType> = non_text.iter().map(|t| **t).collect();
return resolve_union_common_type(&concrete);
}
None
}
fn unify_union_columns(columns: &mut [ColumnSchema], rows: &mut [Row<'static>]) {
for col_idx in 0..columns.len() {
let mut seen: Vec<DataType> = Vec::new();
for row in rows.iter() {
if let Some(dt) = row.values.get(col_idx).and_then(Value::data_type) {
if !seen.contains(&dt) {
seen.push(dt);
}
}
}
if seen.len() == 1
&& matches!(columns[col_idx].ty, DataType::Text)
&& !matches!(seen[0], DataType::Text)
{
columns[col_idx].ty = seen[0];
continue;
}
let Some(target) = resolve_union_common_type(&seen) else {
continue;
};
let scale_preserving_numeric = matches!(target, DataType::Numeric { .. });
let mut coerced: Vec<Option<Value<'static>>> = Vec::with_capacity(rows.len());
let mut ok = true;
for row in rows.iter() {
match row.values.get(col_idx) {
Some(Value::Numeric { .. }) if scale_preserving_numeric => {
coerced.push(Some(row.values[col_idx].clone()));
}
Some(v) => {
let cell_target = if scale_preserving_numeric {
DataType::Numeric {
precision: 0,
scale: 0,
}
} else {
target
};
match crate::conversions::coerce_value(
v.clone(),
cell_target,
&columns[col_idx].name,
col_idx,
) {
Ok(cv) => coerced.push(Some(cv)),
Err(_) => {
ok = false;
break;
}
}
}
None => coerced.push(None),
}
}
if !ok {
continue;
}
for (row, cv) in rows.iter_mut().zip(coerced) {
if let (Some(slot), Some(nv)) = (row.values.get_mut(col_idx), cv) {
*slot = nv;
}
}
columns[col_idx].ty = target;
}
}
fn encode_row_key(row: &Row<'static>) -> Vec<u8> {
let mut out = Vec::new();
for v in &row.values {
match v {
Value::SmallInt(n) => encode_numeric_key(&mut out, i128::from(*n), 0),
Value::Int(n) => encode_numeric_key(&mut out, i128::from(*n), 0),
Value::BigInt(n) => encode_numeric_key(&mut out, i128::from(*n), 0),
Value::Numeric { scaled, scale, .. } => encode_numeric_key(&mut out, *scaled, *scale),
other => {
let s = alloc::format!("{other:?}|");
out.extend_from_slice(s.as_bytes());
}
}
}
out
}
fn encode_numeric_key(out: &mut Vec<u8>, mut scaled: i128, mut scale: u16) {
while scale > 0 && scaled % 10 == 0 {
scaled /= 10;
scale -= 1;
}
let s = alloc::format!("\u{1}{scaled}e-{scale}|");
out.extend_from_slice(s.as_bytes());
}
pub(crate) fn unnest_zip_rows(
args: &[Expr],
) -> Result<(alloc::vec::Vec<DataType>, alloc::vec::Vec<Row<'static>>), EngineError> {
let empty_schema: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let ctx = EvalContext::new(&empty_schema, None);
let dummy_row = Row::new(alloc::vec::Vec::new());
let mut dtypes: alloc::vec::Vec<DataType> = alloc::vec::Vec::with_capacity(args.len());
let mut columns: alloc::vec::Vec<alloc::vec::Vec<Value<'static>>> =
alloc::vec::Vec::with_capacity(args.len());
for a in args {
let v = eval::eval_expr(a, &dummy_row, &ctx).map_err(EngineError::Eval)?;
let (dt, items): (DataType, alloc::vec::Vec<Value<'static>>) = match v {
Value::Null => (DataType::Text, alloc::vec::Vec::new()),
Value::TextArray(xs) => (
DataType::Text,
xs.into_iter()
.map(|x| x.map(Value::text).unwrap_or(Value::Null))
.collect(),
),
Value::IntArray(xs) => (
DataType::Int,
xs.into_iter()
.map(|x| x.map(Value::Int).unwrap_or(Value::Null))
.collect(),
),
Value::BigIntArray(xs) => (
DataType::BigInt,
xs.into_iter()
.map(|x| x.map(Value::BigInt).unwrap_or(Value::Null))
.collect(),
),
other => {
return Err(EngineError::Unsupported(alloc::format!(
"unnest() expects array arguments, got {}",
crate::conversions::pg_type_name_for_error_opt(other.data_type())
)));
}
};
dtypes.push(dt);
columns.push(items);
}
let max_len = columns.iter().map(|c| c.len()).max().unwrap_or(0);
let mut rows: alloc::vec::Vec<Row<'static>> = alloc::vec::Vec::with_capacity(max_len);
for i in 0..max_len {
let vals: alloc::vec::Vec<Value<'static>> = columns
.iter()
.map(|c| c.get(i).cloned().unwrap_or(Value::Null))
.collect();
rows.push(Row::new(vals));
}
Ok((dtypes, rows))
}
pub(crate) fn unnest_zip_args(expr: &Expr) -> Option<&[Expr]> {
match expr {
Expr::FunctionCall { name, args } if name == "__unnest_zip" => Some(args.as_slice()),
_ => None,
}
}
pub(crate) fn generate_series_rows(
args: &[Expr],
cancel: &CancelToken<'_>,
) -> Result<(DataType, alloc::vec::Vec<Row<'static>>), EngineError> {
let empty_schema: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let ctx = EvalContext::new(&empty_schema, None);
let dummy_row = Row::new(alloc::vec::Vec::new());
let mut arg_values: alloc::vec::Vec<Value<'static>> =
alloc::vec::Vec::with_capacity(args.len());
for a in args {
arg_values.push(eval::eval_expr(a, &dummy_row, &ctx).map_err(EngineError::Eval)?);
}
generate_series_from_values(arg_values, args, cancel)
}
pub(crate) fn generate_series_from_values(
mut arg_values: alloc::vec::Vec<Value<'static>>,
args: &[Expr],
cancel: &CancelToken<'_>,
) -> Result<(DataType, alloc::vec::Vec<Row<'static>>), EngineError> {
if arg_values.iter().any(|v| matches!(v, Value::Null)) {
return Ok((DataType::BigInt, alloc::vec::Vec::new()));
}
let empty_cols: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let tz = arg_values.iter().any(|v| matches!(v, Value::Date(_)))
|| args.iter().any(|a| {
crate::describe::describe_expr(a, &empty_cols)
.is_some_and(|s| matches!(s.ty, DataType::Timestamptz))
});
for v in &mut arg_values {
if let Value::Date(d) = *v {
*v = Value::Timestamp(crate::conversions::date_days_to_micros(d));
}
}
match arg_values.as_slice() {
[Value::Timestamp(start), Value::Timestamp(stop), step] => {
let interval_step = match step {
Value::Interval { .. } => step.clone(),
Value::Text(s) => crate::conversions::coerce_value(
Value::text(s.as_ref()),
DataType::Interval,
"",
0,
)
.map_err(|_| {
EngineError::Unsupported(alloc::format!(
"generate_series(timestamp, timestamp, …): \
could not parse step {s:?} as INTERVAL"
))
})?,
other => {
return Err(EngineError::Unsupported(alloc::format!(
"generate_series(timestamp, timestamp, …): \
step must be INTERVAL, got {}",
crate::conversions::pg_type_name_for_error_opt(other.data_type())
)));
}
};
let rows = generate_series_timestamps(*start, *stop, interval_step, cancel)?;
Ok((
if tz {
DataType::Timestamptz
} else {
DataType::Timestamp
},
rows,
))
}
[start, stop, step]
if value_is_integer(start) && value_is_integer(stop) && value_is_integer(step) =>
{
let s = value_to_i64(start);
let e = value_to_i64(stop);
let st = value_to_i64(step);
let wide = value_is_bigint(start) || value_is_bigint(stop) || value_is_bigint(step);
let rows = generate_series_integers(s, e, st, wide, cancel)?;
Ok((
if wide {
DataType::BigInt
} else {
DataType::Int
},
rows,
))
}
[start, stop] if value_is_integer(start) && value_is_integer(stop) => {
let s = value_to_i64(start);
let e = value_to_i64(stop);
let wide = value_is_bigint(start) || value_is_bigint(stop);
let rows = generate_series_integers(s, e, 1, wide, cancel)?;
Ok((
if wide {
DataType::BigInt
} else {
DataType::Int
},
rows,
))
}
[_, _] | [_, _, _]
if arg_values
.iter()
.any(|v| matches!(v, Value::Numeric { .. } | Value::NumericBig(_)))
&& arg_values.iter().all(|v| {
matches!(v, Value::Numeric { .. } | Value::NumericBig(_)) || value_is_integer(v)
}) =>
{
use spg_storage::NumericKind as K;
let words: [(&str, &str); 3] = [
(
"start value cannot be NaN",
"start value cannot be infinity",
),
("stop value cannot be NaN", "stop value cannot be infinity"),
("step size cannot be NaN", "step size cannot be infinity"),
];
for (i, v) in arg_values.iter().enumerate() {
if let Value::Numeric { kind, .. } = v {
if *kind != K::Finite {
let (nan_w, inf_w) = words[i];
return Err(EngineError::Unsupported(
if *kind == K::NaN { nan_w } else { inf_w }.into(),
));
}
}
}
let big =
|v: &Value<'_>| eval::binop::value_to_bignum(v).expect("finite numeric or integer");
let start = big(&arg_values[0]);
let stop = big(&arg_values[1]);
let step = if arg_values.len() == 3 {
big(&arg_values[2])
} else {
spg_storage::bignum::BigNumeric::from_i128(1, 0)
};
if step.is_zero() {
return Err(EngineError::Unsupported(
"step size cannot equal zero".into(),
));
}
let descending = step.parts().0;
let mut rows = alloc::vec::Vec::new();
let mut cur = start;
const MAX_ROWS: usize = 10_000_000;
loop {
cancel.check()?;
let c = cur.cmp(&stop);
if descending {
if c == core::cmp::Ordering::Less {
break;
}
} else if c == core::cmp::Ordering::Greater {
break;
}
if rows.len() >= MAX_ROWS {
return Err(EngineError::Unsupported(alloc::format!(
"generate_series() result exceeds {MAX_ROWS} rows"
)));
}
rows.push(Row::new(alloc::vec![eval::binop::bignum_to_value(
cur.clone()
)]));
cur = cur.add(&step);
}
Ok((
DataType::Numeric {
precision: 0,
scale: 0,
},
rows,
))
}
_ => Err(EngineError::Unsupported(alloc::format!(
"generate_series(): v7.17 supports integer or (timestamp, timestamp, interval) \
argument shapes; got {}",
arg_values
.iter()
.map(|v| crate::conversions::pg_type_name_for_error_opt(v.data_type()))
.collect::<alloc::vec::Vec<_>>()
.join(", ")
))),
}
}
fn generate_series_integers(
start: i64,
stop: i64,
step: i64,
wide: bool,
cancel: &CancelToken<'_>,
) -> Result<alloc::vec::Vec<Row<'static>>, EngineError> {
if step == 0 {
return Err(EngineError::Unsupported(
"step size cannot equal zero".into(),
));
}
let mut out = alloc::vec::Vec::new();
let mut cur = start;
const MAX_ROWS: usize = 10_000_000;
loop {
cancel.check()?;
if step > 0 && cur > stop {
break;
}
if step < 0 && cur < stop {
break;
}
out.push(Row::new(alloc::vec![if wide {
Value::BigInt(cur)
} else {
Value::Int(cur as i32)
}]));
if out.len() > MAX_ROWS {
return Err(EngineError::Unsupported(alloc::format!(
"generate_series(): exceeded {MAX_ROWS} rows; \
narrow start/stop or use a larger step"
)));
}
cur = match cur.checked_add(step) {
Some(n) => n,
None => break,
};
}
Ok(out)
}
fn generate_series_timestamps(
start: i64,
stop: i64,
step: Value,
cancel: &CancelToken<'_>,
) -> Result<alloc::vec::Vec<Row<'static>>, EngineError> {
let (months, days, micros) = match &step {
Value::Interval {
months,
days,
micros,
} => (*months, *days, *micros),
_ => unreachable!("caller guards step.is_interval"),
};
if months == 0 && days == 0 && micros == 0 {
return Err(EngineError::Unsupported(
"generate_series(): INTERVAL step cannot be zero".into(),
));
}
let ascending = months > 0 || days > 0 || micros > 0;
let mut out = alloc::vec::Vec::new();
let mut cur = Value::Timestamp(start);
const MAX_ROWS: usize = 10_000_000;
loop {
cancel.check()?;
let cur_t = match cur {
Value::Timestamp(t) => t,
_ => unreachable!("loop invariant: cur is Timestamp"),
};
if ascending && cur_t > stop {
break;
}
if !ascending && cur_t < stop {
break;
}
out.push(Row::new(alloc::vec![Value::Timestamp(cur_t)]));
if out.len() > MAX_ROWS {
return Err(EngineError::Unsupported(alloc::format!(
"generate_series(): exceeded {MAX_ROWS} rows; \
narrow start/stop or use a larger step"
)));
}
let next = eval::apply_binary_interval(
spg_sql::ast::BinOp::Add,
&cur,
&Value::Interval {
months,
days,
micros,
},
)
.map_err(EngineError::Eval)?;
cur = match next {
Some(v) => v,
None => break,
};
}
Ok(out)
}
fn check_with_ties_requires_order_by(stmt: &SelectStatement) -> Result<(), EngineError> {
if stmt.limit_with_ties && stmt.order_by.is_empty() {
return Err(EngineError::Unsupported(alloc::string::String::from(
"WITH TIES cannot be specified without ORDER BY clause",
)));
}
Ok(())
}
fn is_top_level_unnest(expr: &spg_sql::ast::Expr) -> bool {
top_level_srf_kind(expr).is_some()
}
#[derive(Clone, Copy, PartialEq, Eq)]
pub(crate) enum SrfKind {
Unnest,
GenerateSeries,
GenerateSubscripts,
ArrayElements {
as_text: bool,
},
PathQuery,
RegexpMatches,
Each {
as_text: bool,
},
ObjectKeys,
}
fn name_is(name: &str, names: &[&str]) -> bool {
names.iter().any(|n| name.eq_ignore_ascii_case(n))
}
pub(crate) fn top_level_srf_kind(expr: &spg_sql::ast::Expr) -> Option<SrfKind> {
let spg_sql::ast::Expr::FunctionCall { name, args } = expr else {
return None;
};
let n = args.len();
if n == 1 && name.eq_ignore_ascii_case("unnest") {
return Some(SrfKind::Unnest);
}
if (2..=3).contains(&n) && name.eq_ignore_ascii_case("generate_series") {
return Some(SrfKind::GenerateSeries);
}
if n == 2 && name.eq_ignore_ascii_case("generate_subscripts") {
return Some(SrfKind::GenerateSubscripts);
}
if n == 1 && name_is(name, &["jsonb_array_elements", "json_array_elements"]) {
return Some(SrfKind::ArrayElements { as_text: false });
}
if n == 1
&& name_is(
name,
&["jsonb_array_elements_text", "json_array_elements_text"],
)
{
return Some(SrfKind::ArrayElements { as_text: true });
}
if (2..=4).contains(&n) && name_is(name, &["jsonb_path_query", "json_path_query"]) {
return Some(SrfKind::PathQuery);
}
if (2..=3).contains(&n) && name.eq_ignore_ascii_case("regexp_matches") {
return Some(SrfKind::RegexpMatches);
}
if n == 1 && name_is(name, &["jsonb_each", "json_each"]) {
return Some(SrfKind::Each { as_text: false });
}
if n == 1 && name_is(name, &["jsonb_each_text", "json_each_text"]) {
return Some(SrfKind::Each { as_text: true });
}
if n == 1 && name_is(name, &["jsonb_object_keys", "json_object_keys"]) {
return Some(SrfKind::ObjectKeys);
}
None
}
pub(crate) fn top_level_srf_output(
expr: &spg_sql::ast::Expr,
row: &Row<'static>,
ctx: &EvalContext<'_>,
) -> Result<Vec<Value<'static>>, EngineError> {
let (Some(kind), spg_sql::ast::Expr::FunctionCall { name, args }) =
(top_level_srf_kind(expr), expr)
else {
return Err(EngineError::Unsupported(
"expected a SELECT-list SRF call".into(),
));
};
match kind {
SrfKind::Unnest => {
if let spg_sql::ast::Expr::Array(items) = &args[0] {
return items
.iter()
.map(|e| eval::eval_expr(e, row, ctx).map_err(EngineError::Eval))
.collect();
}
let arr = eval::eval_expr(&args[0], row, ctx).map_err(EngineError::Eval)?;
array_value_to_elements(&arr)
}
SrfKind::GenerateSeries => {
let mut arg_values: Vec<Value<'static>> = Vec::with_capacity(args.len());
for a in args {
arg_values.push(eval::eval_expr(a, row, ctx).map_err(EngineError::Eval)?);
}
let (_, rows) = generate_series_from_values(arg_values, args, &CancelToken::none())?;
Ok(rows
.into_iter()
.map(|r| r.values.into_iter().next().unwrap_or(Value::Null))
.collect())
}
SrfKind::GenerateSubscripts => {
let arr = eval::eval_expr(&args[0], row, ctx).map_err(EngineError::Eval)?;
let dim = eval::eval_expr(&args[1], row, ctx).map_err(EngineError::Eval)?;
if !matches!(dim, Value::Int(1) | Value::BigInt(1) | Value::SmallInt(1)) {
return Ok(Vec::new());
}
let len = array_value_to_elements(&arr)?.len();
Ok((1..=len).map(|i| Value::Int(i as i32)).collect())
}
SrfKind::ArrayElements { as_text } => {
let arg = eval::eval_expr(&args[0], row, ctx).map_err(EngineError::Eval)?;
if matches!(arg, Value::Null) {
return Ok(Vec::new());
}
let items =
crate::json::array_element_rows(&arg, as_text, name).map_err(EngineError::Eval)?;
Ok(items
.into_iter()
.map(|opt| opt.map(Value::text).unwrap_or(Value::Null))
.collect())
}
SrfKind::ObjectKeys => {
let v = eval::eval_expr(expr, row, ctx).map_err(EngineError::Eval)?;
array_value_to_elements(&v)
}
SrfKind::RegexpMatches => {
let vals: Vec<Value<'static>> = args
.iter()
.map(|a| eval::eval_expr(a, row, ctx).map_err(EngineError::Eval))
.collect::<Result<_, _>>()?;
crate::eval::regexp_matches_rows(&vals).map_err(EngineError::Eval)
}
SrfKind::Each { as_text } => {
let arg = eval::eval_expr(&args[0], row, ctx).map_err(EngineError::Eval)?;
if matches!(arg, Value::Null) {
return Ok(Vec::new());
}
let pairs = crate::json::each_rows(&arg, as_text, name).map_err(EngineError::Eval)?;
Ok(pairs
.into_iter()
.map(|(k, v)| {
let val = if as_text {
v.map(Value::text).unwrap_or(Value::Null)
} else {
v.map(Value::json).unwrap_or(Value::Null)
};
Value::Composite(alloc::vec![
("key".to_string(), Value::text(k)),
("value".to_string(), val),
])
})
.collect())
}
SrfKind::PathQuery => {
let doc = eval::eval_expr(&args[0], row, ctx).map_err(EngineError::Eval)?;
let path = eval::eval_expr(&args[1], row, ctx).map_err(EngineError::Eval)?;
let vars = match args.get(2) {
Some(a) => {
let v = eval::eval_expr(a, row, ctx).map_err(EngineError::Eval)?;
crate::json::parse_path_vars(&v).map_err(EngineError::Eval)?
}
None => None,
};
match crate::json::path_query_vars(&doc, &path, vars.as_ref())
.map_err(EngineError::Eval)?
{
Value::Null => Ok(Vec::new()),
Value::TextArray(items) => Ok(items
.into_iter()
.map(|opt| opt.map(Value::text).unwrap_or(Value::Null))
.collect()),
other => Ok(alloc::vec![other]),
}
}
}
}
pub(crate) fn array_value_to_elements(v: &Value) -> Result<Vec<Value<'static>>, EngineError> {
if let Some(flat) = crate::eval::values::flatten_2d(v) {
return array_value_to_elements(&flat);
}
match v {
Value::Null => Ok(Vec::new()),
Value::TextArray(items) => Ok(items
.iter()
.map(|opt| {
opt.as_ref()
.map(|s| Value::text(s.clone()))
.unwrap_or(Value::Null)
})
.collect()),
Value::IntArray(items) => Ok(items
.iter()
.map(|opt| opt.map(Value::Int).unwrap_or(Value::Null))
.collect()),
Value::BigIntArray(items) => Ok(items
.iter()
.map(|opt| opt.map(Value::BigInt).unwrap_or(Value::Null))
.collect()),
Value::Multirange { kind, ranges } => Ok(ranges
.iter()
.map(|s| Value::Range {
kind: *kind,
lower: s.lower.clone(),
upper: s.upper.clone(),
lower_inc: s.lower_inc,
upper_inc: s.upper_inc,
empty: false,
})
.collect()),
other => Err(EngineError::Eval(EvalError::TypeMismatch {
detail: alloc::format!(
"unnest() expects an array argument, got {}",
crate::conversions::pg_type_name_for_error_opt(other.data_type())
),
})),
}
}
impl Engine {
fn expand_views_in_select(
&self,
stmt: &SelectStatement,
) -> Result<Option<SelectStatement>, EngineError> {
let cat = self.active_catalog();
let mut referenced: Vec<String> = Vec::new();
if let Some(from) = &stmt.from {
collect_view_refs(&from.primary, cat, &mut referenced);
for j in &from.joins {
collect_view_refs(&j.table, cat, &mut referenced);
}
}
referenced.retain(|n| !stmt.ctes.iter().any(|c| c.name == *n));
if referenced.is_empty() {
return Ok(None);
}
let mut new_ctes: Vec<spg_sql::ast::Cte> = Vec::with_capacity(referenced.len());
for name in &referenced {
let view = cat.view(name).ok_or_else(|| {
EngineError::Storage(spg_storage::StorageError::Corrupt(alloc::format!(
"view {name:?} disappeared mid-expansion"
)))
})?;
let parsed = spg_sql::parser::parse_statement(&view.body).map_err(|e| {
EngineError::Unsupported(alloc::format!("view {name:?} body re-parse failed: {e}"))
})?;
let Statement::Select(body) = parsed else {
return Err(EngineError::Unsupported(alloc::format!(
"view {name:?} body is not a SELECT (catalog corruption)"
)));
};
new_ctes.push(spg_sql::ast::Cte {
name: name.clone(),
body: spg_sql::ast::CteBody::Select(body),
recursive: false,
column_overrides: view.columns.clone(),
search: None,
cycle: None,
});
}
let mut out = stmt.clone();
new_ctes.extend(out.ctes);
out.ctes = new_ctes;
Ok(Some(out))
}
fn expand_partition_parents_in_select(
&self,
stmt: &SelectStatement,
) -> Result<Option<SelectStatement>, EngineError> {
let cat = self.active_catalog();
let Some(from) = &stmt.from else {
return Ok(None);
};
let mut parent_refs: Vec<String> = Vec::new();
collect_partition_parent_refs(&from.primary, cat, &mut parent_refs);
for j in &from.joins {
collect_partition_parent_refs(&j.table, cat, &mut parent_refs);
}
parent_refs.retain(|n| !stmt.ctes.iter().any(|c| c.name.eq_ignore_ascii_case(n)));
if parent_refs.is_empty() {
return Ok(None);
}
let synth_name = |p: &str| alloc::format!("__spg_partition_{p}");
let mut new_ctes: Vec<spg_sql::ast::Cte> = Vec::with_capacity(parent_refs.len());
let mut expanded_parents: Vec<alloc::string::String> = Vec::new();
for parent_name in &parent_refs {
let Some(body) = self.build_partition_parent_union_body(parent_name, stmt)? else {
continue;
};
new_ctes.push(spg_sql::ast::Cte {
name: synth_name(parent_name),
body: spg_sql::ast::CteBody::Select(body),
recursive: false,
column_overrides: Vec::new(),
search: None,
cycle: None,
});
expanded_parents.push(parent_name.clone());
}
if expanded_parents.is_empty() {
return Ok(None);
}
let mut out = stmt.clone();
if let Some(from) = out.from.as_mut() {
rewrite_partition_parent_table_ref(&mut from.primary, &expanded_parents, &synth_name);
for j in &mut from.joins {
rewrite_partition_parent_table_ref(&mut j.table, &expanded_parents, &synth_name);
}
}
new_ctes.extend(out.ctes);
out.ctes = new_ctes;
Ok(Some(out))
}
pub(crate) fn explain_partition_kept_children_by_where(
&self,
parent_name: &str,
where_: Option<&spg_sql::ast::Expr>,
) -> Option<Vec<alloc::string::String>> {
let mut synth = SelectStatement::default();
synth.where_ = where_.cloned();
self.explain_partition_kept_children(parent_name, &synth)
}
pub(crate) fn explain_partition_kept_children(
&self,
parent_name: &str,
outer: &SelectStatement,
) -> Option<Vec<alloc::string::String>> {
use spg_storage::PartitionRole;
let cat = self.active_catalog();
let parent = cat.get(parent_name)?;
let (key_position, parent_kind) = match &parent.schema().partition_role {
Some(PartitionRole::Parent {
key_column_positions,
kind,
..
}) => (*key_column_positions.first().unwrap_or(&0), *kind),
_ => return None,
};
let key_col_name = parent.schema().columns[key_position].name.clone();
let (lo_bound, hi_bound) = match outer.where_.as_ref() {
Some(expr) => extract_key_range(expr, &key_col_name),
None => (None, None),
};
let eq_value: Option<spg_storage::Value<'static>> = match outer.where_.as_ref() {
Some(expr) => extract_key_eq_value(expr, &key_col_name),
None => None,
};
let children = crate::partition::children_of_parent(cat, parent_name);
let mut kept: Vec<alloc::string::String> = Vec::new();
let mut default_child: Option<alloc::string::String> = None;
for child_name in &children {
let Some(child) = cat.get(child_name) else {
continue;
};
match &child.schema().partition_role {
Some(PartitionRole::Range { lower, upper, .. }) => {
if range_satisfies_filter(lower, upper, lo_bound.as_ref(), hi_bound.as_ref()) {
kept.push(child_name.clone());
}
}
Some(PartitionRole::List { values, .. }) => match &eq_value {
Some(v) => {
if values.iter().any(|b| b.equals_value(v)) {
kept.push(child_name.clone());
}
}
None => kept.push(child_name.clone()),
},
Some(PartitionRole::Hash {
modulus, remainder, ..
}) => match &eq_value {
Some(v) => {
let h = crate::partition::pg_compatible_hash(v);
if h.rem_euclid(u64::from(*modulus)) == u64::from(*remainder) {
kept.push(child_name.clone());
}
}
None => kept.push(child_name.clone()),
},
Some(PartitionRole::Default { .. }) => {
default_child = Some(child_name.clone());
}
_ => {}
}
}
let _ = parent_kind;
if let Some(d) = default_child {
if kept.is_empty() || eq_value.is_none() {
kept.push(d);
}
}
Some(kept)
}
fn build_partition_parent_union_body(
&self,
parent_name: &str,
outer: &SelectStatement,
) -> Result<Option<SelectStatement>, EngineError> {
use spg_storage::PartitionRole;
let cat = self.active_catalog();
let parent = cat.get(parent_name).ok_or_else(|| {
EngineError::Storage(spg_storage::StorageError::Corrupt(alloc::format!(
"partition parent {parent_name:?} disappeared mid-expansion"
)))
})?;
let (key_position, parent_kind) = match &parent.schema().partition_role {
Some(PartitionRole::Parent {
key_column_positions,
kind,
..
}) => (*key_column_positions.first().unwrap_or(&0), *kind),
_ if crate::partition::has_inheritance_children(cat, parent_name) => {
let cols = parent
.schema()
.columns
.iter()
.map(|c| quote_ident_for_sql(&c.name))
.collect::<Vec<_>>()
.join(", ");
let carry_sys = references_ctid(outer);
let sys = if carry_sys {
let mut t = alloc::string::String::new();
for s in SYSTEM_COLUMNS {
t.push_str(", ");
t.push_str(s);
}
t
} else {
alloc::string::String::new()
};
let mut body = alloc::format!(
"SELECT {cols}{sys} FROM ONLY {}",
quote_ident_for_sql(parent_name)
);
for child in crate::partition::children_of_parent(cat, parent_name) {
body.push_str(&alloc::format!(
" UNION ALL SELECT {cols}{sys} FROM {}",
quote_ident_for_sql(&child)
));
}
return parse_select_or_corrupt(&body).map(Some);
}
_ => {
return Err(EngineError::Unsupported(alloc::format!(
"partition expansion: {parent_name:?} is not a parent"
)));
}
};
let key_col_name = parent.schema().columns[key_position].name.clone();
let (lo_bound, hi_bound) = match outer.where_.as_ref() {
Some(expr) => extract_key_range(expr, &key_col_name),
None => (None, None),
};
let eq_value: Option<spg_storage::Value<'static>> = match outer.where_.as_ref() {
Some(expr) => extract_key_eq_value(expr, &key_col_name),
None => None,
};
let children = crate::partition::children_of_parent(cat, parent_name);
let mut kept: Vec<String> = Vec::new();
let mut default_child: Option<String> = None;
for child_name in &children {
let Some(child) = cat.get(child_name) else {
continue;
};
match &child.schema().partition_role {
Some(PartitionRole::Range { lower, upper, .. }) => {
if range_satisfies_filter(lower, upper, lo_bound.as_ref(), hi_bound.as_ref()) {
kept.push(child_name.clone());
}
}
Some(PartitionRole::List { values, .. }) => match &eq_value {
Some(v) => {
if values.iter().any(|b| b.equals_value(v)) {
kept.push(child_name.clone());
}
}
None => kept.push(child_name.clone()),
},
Some(PartitionRole::Hash {
modulus, remainder, ..
}) => match &eq_value {
Some(v) => {
let h = crate::partition::pg_compatible_hash(v);
if h.rem_euclid(u64::from(*modulus)) == u64::from(*remainder) {
kept.push(child_name.clone());
}
}
None => kept.push(child_name.clone()),
},
Some(PartitionRole::Default { .. }) => {
default_child = Some(child_name.clone());
}
_ => {}
}
}
let _ = parent_kind; if let Some(d) = default_child {
if kept.is_empty() {
kept.push(d);
} else if eq_value.is_none() {
kept.push(d);
}
}
if kept.is_empty() {
let _ = parent_name;
return Ok(None);
}
let carry_sys = references_ctid(outer);
let mut body = alloc::string::String::new();
for (i, child_name) in kept.iter().enumerate() {
if i > 0 {
body.push_str(" UNION ALL ");
}
body.push_str("SELECT *");
if carry_sys {
for sys in SYSTEM_COLUMNS {
body.push_str(", ");
body.push_str(sys);
}
}
body.push_str(" FROM ");
body.push_str("e_ident_for_sql(child_name));
}
parse_select_or_corrupt(&body).map(Some)
}
}
fn rewrite_partition_parent_table_ref(
t: &mut spg_sql::ast::TableRef,
parents: &[alloc::string::String],
synth_name: &impl Fn(&str) -> alloc::string::String,
) {
if t.lateral_subquery.is_some() || t.unnest_expr.is_some() || t.generate_series_args.is_some() {
return;
}
if t.only || !parents.iter().any(|p| p == &t.name) {
return;
}
if t.alias.is_none() {
t.alias = Some(t.name.clone());
}
t.name = synth_name(&t.name);
}
fn collect_partition_parent_refs(
t: &spg_sql::ast::TableRef,
cat: &spg_storage::Catalog,
out: &mut Vec<alloc::string::String>,
) {
if t.lateral_subquery.is_some() || t.unnest_expr.is_some() || t.generate_series_args.is_some() {
return;
}
if !t.only && crate::partition::has_children(cat, &t.name) {
out.push(t.name.clone());
}
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct PartitionFilterBound {
pub micros: i64,
pub inclusive: bool,
}
fn extract_key_range(
expr: &spg_sql::ast::Expr,
key_col: &str,
) -> (Option<PartitionFilterBound>, Option<PartitionFilterBound>) {
let mut lo: Option<PartitionFilterBound> = None;
let mut hi: Option<PartitionFilterBound> = None;
let mut stack: Vec<&spg_sql::ast::Expr> = alloc::vec![expr];
while let Some(e) = stack.pop() {
match e {
spg_sql::ast::Expr::Binary {
lhs,
op: spg_sql::ast::BinOp::And,
rhs,
} => {
stack.push(lhs);
stack.push(rhs);
}
spg_sql::ast::Expr::Binary { lhs, op, rhs } => {
let (col_ref, lit_side, swapped) = if is_column_ref(lhs, key_col) {
(Some(lhs.as_ref()), rhs.as_ref(), false)
} else if is_column_ref(rhs, key_col) {
(Some(rhs.as_ref()), lhs.as_ref(), true)
} else {
(None, lhs.as_ref(), false)
};
if col_ref.is_none() {
continue;
}
let Some(lit) = literal_to_micros(lit_side) else {
continue;
};
use spg_sql::ast::BinOp::{Eq, Gt, GtEq, Lt, LtEq};
let effective_op = if swapped {
match op {
Lt => Gt,
LtEq => GtEq,
Gt => Lt,
GtEq => LtEq,
other => *other,
}
} else {
*op
};
match effective_op {
Eq => {
tighten_lo(
&mut lo,
PartitionFilterBound {
micros: lit,
inclusive: true,
},
);
tighten_hi(
&mut hi,
PartitionFilterBound {
micros: lit,
inclusive: true,
},
);
}
GtEq => {
tighten_lo(
&mut lo,
PartitionFilterBound {
micros: lit,
inclusive: true,
},
);
}
Gt => {
tighten_lo(
&mut lo,
PartitionFilterBound {
micros: lit,
inclusive: false,
},
);
}
LtEq => {
tighten_hi(
&mut hi,
PartitionFilterBound {
micros: lit,
inclusive: true,
},
);
}
Lt => {
tighten_hi(
&mut hi,
PartitionFilterBound {
micros: lit,
inclusive: false,
},
);
}
_ => {}
}
}
_ => {}
}
}
(lo, hi)
}
fn tighten_lo(slot: &mut Option<PartitionFilterBound>, new: PartitionFilterBound) {
match slot {
None => *slot = Some(new),
Some(cur) => {
if new.micros > cur.micros
|| (new.micros == cur.micros && !new.inclusive && cur.inclusive)
{
*slot = Some(new);
}
}
}
}
fn tighten_hi(slot: &mut Option<PartitionFilterBound>, new: PartitionFilterBound) {
match slot {
None => *slot = Some(new),
Some(cur) => {
if new.micros < cur.micros
|| (new.micros == cur.micros && !new.inclusive && cur.inclusive)
{
*slot = Some(new);
}
}
}
}
fn is_column_ref(e: &spg_sql::ast::Expr, key_col: &str) -> bool {
if let spg_sql::ast::Expr::Column(c) = e {
c.name.eq_ignore_ascii_case(key_col)
} else {
false
}
}
pub(crate) fn extract_key_eq_value(
expr: &spg_sql::ast::Expr,
key_col: &str,
) -> Option<spg_storage::Value<'static>> {
let mut stack: Vec<&spg_sql::ast::Expr> = alloc::vec![expr];
while let Some(e) = stack.pop() {
match e {
spg_sql::ast::Expr::Binary {
lhs,
op: spg_sql::ast::BinOp::And,
rhs,
} => {
stack.push(lhs);
stack.push(rhs);
}
spg_sql::ast::Expr::Binary {
lhs,
op: spg_sql::ast::BinOp::Eq,
rhs,
} => {
let lit_side = if is_column_ref(lhs, key_col) {
rhs.as_ref()
} else if is_column_ref(rhs, key_col) {
lhs.as_ref()
} else {
continue;
};
let cloned = lit_side.clone();
let Ok(v) = crate::conversions::literal_expr_to_value(cloned) else {
continue;
};
let owned: spg_storage::Value<'static> = match v {
spg_storage::Value::Text(s) => {
spg_storage::Value::Text(alloc::borrow::Cow::Owned(s.into_owned()))
}
spg_storage::Value::SmallInt(n) => spg_storage::Value::SmallInt(n),
spg_storage::Value::Int(n) => spg_storage::Value::Int(n),
spg_storage::Value::BigInt(n) => spg_storage::Value::BigInt(n),
spg_storage::Value::Date(d) => spg_storage::Value::Date(d),
spg_storage::Value::Timestamp(t) => spg_storage::Value::Timestamp(t),
spg_storage::Value::Bool(b) => spg_storage::Value::Bool(b),
spg_storage::Value::Null => spg_storage::Value::Null,
_ => continue,
};
return Some(owned);
}
_ => {}
}
}
None
}
fn literal_to_micros(e: &spg_sql::ast::Expr) -> Option<i64> {
let cloned = e.clone();
let value = crate::conversions::literal_expr_to_value(cloned).ok()?;
match value {
spg_storage::Value::Timestamp(m) => Some(m),
spg_storage::Value::Date(days) => Some(i64::from(days) * 86_400i64 * 1_000_000i64),
spg_storage::Value::Text(s) => crate::eval::parse_timestamp_literal(&s),
_ => None,
}
}
fn range_satisfies_filter(
range_lo: &spg_storage::PartitionBound,
range_hi: &spg_storage::PartitionBound,
filter_lo: Option<&PartitionFilterBound>,
filter_hi: Option<&PartitionFilterBound>,
) -> bool {
use spg_storage::PartitionBound;
if let Some(lo) = filter_lo {
match range_hi {
PartitionBound::MinValue => return false,
PartitionBound::MaxValue => {}
PartitionBound::TimestampTz(hi) => {
if *hi <= lo.micros {
return false;
}
}
PartitionBound::BigInt(_)
| PartitionBound::Int(_)
| PartitionBound::SmallInt(_)
| PartitionBound::Date(_)
| PartitionBound::Text(_) => {}
}
}
if let Some(hi) = filter_hi {
match range_lo {
PartitionBound::MaxValue => return false,
PartitionBound::MinValue => {}
PartitionBound::TimestampTz(lo) => {
let rejects = if hi.inclusive {
*lo > hi.micros
} else {
*lo >= hi.micros
};
if rejects {
return false;
}
}
PartitionBound::BigInt(_)
| PartitionBound::Int(_)
| PartitionBound::SmallInt(_)
| PartitionBound::Date(_)
| PartitionBound::Text(_) => {}
}
}
true
}
fn quote_ident_for_sql(name: &str) -> alloc::string::String {
let mut out = alloc::string::String::with_capacity(name.len() + 2);
out.push('"');
for c in name.chars() {
if c == '"' {
out.push('"');
}
out.push(c);
}
out.push('"');
out
}
fn parse_select_or_corrupt(sql: &str) -> Result<SelectStatement, EngineError> {
let parsed = spg_sql::parser::parse_statement(sql).map_err(|e| {
EngineError::Unsupported(alloc::format!(
"partition expansion: generated SQL {sql:?} failed to re-parse: {e}"
))
})?;
let Statement::Select(body) = parsed else {
return Err(EngineError::Unsupported(alloc::format!(
"partition expansion: generated SQL {sql:?} is not a SELECT"
)));
};
Ok(body)
}
fn setof_column_shape_from(
declared: &str,
name: &str,
alias: Option<&str>,
got: &[ColumnSchema],
) -> alloc::vec::Vec<ColumnSchema> {
let upper = declared.to_ascii_uppercase();
if upper.starts_with("TABLE(") {
let raw = &declared["TABLE(".len()..declared.len() - 1];
return raw
.split(',')
.zip(got.iter())
.map(|(decl, g)| {
let cname = decl.split_whitespace().next().unwrap_or(g.name.as_str());
ColumnSchema::new(cname.to_string(), g.ty, true)
})
.collect();
}
let cname = alias.unwrap_or(name);
got.first()
.map(|c| alloc::vec![ColumnSchema::new(cname.to_string(), c.ty, true)])
.unwrap_or_default()
}
fn setof_column_shape(
declared: &str,
name: &str,
alias: Option<&str>,
first_row: Option<&alloc::vec::Vec<Value<'static>>>,
) -> alloc::vec::Vec<ColumnSchema> {
let got: alloc::vec::Vec<ColumnSchema> = first_row
.map(|r| {
r.iter()
.enumerate()
.map(|(i, v)| {
ColumnSchema::new(
alloc::format!("col{i}"),
v.data_type().unwrap_or(DataType::Text),
true,
)
})
.collect()
})
.unwrap_or_default();
setof_column_shape_from(declared, name, alias, &got)
}
fn validate_locking_clause(stmt: &SelectStatement) -> Result<(), EngineError> {
let Some(lock) = &stmt.locking else {
return Ok(());
};
let verb = lock_clause_verb(lock.strength);
let refuse = |what: &str| {
Err(EngineError::Unsupported(alloc::format!(
"{verb} is not allowed with {what}"
)))
};
if !stmt.unions.is_empty() {
return refuse("UNION/INTERSECT/EXCEPT");
}
if stmt.distinct || !stmt.distinct_on.is_empty() {
return refuse("DISTINCT clause");
}
if stmt.group_by.is_some() || stmt.group_by_all {
return refuse("GROUP BY clause");
}
let has_agg = stmt.items.iter().any(|it| match it {
spg_sql::ast::SelectItem::Expr { expr, .. } => crate::aggregate::contains_aggregate(expr),
_ => false,
});
if has_agg {
return refuse("aggregate functions");
}
for want in &lock.of_tables {
if !locking_from_names(stmt)
.iter()
.any(|n| n.eq_ignore_ascii_case(want))
{
return Err(EngineError::Unsupported(alloc::format!(
"relation \"{want}\" in {verb} clause not found in FROM clause"
)));
}
}
Ok(())
}
const fn lock_clause_verb(s: spg_sql::ast::LockStrength) -> &'static str {
use spg_sql::ast::LockStrength as LS;
match s {
LS::Update => "FOR UPDATE",
LS::NoKeyUpdate => "FOR NO KEY UPDATE",
LS::Share => "FOR SHARE",
LS::KeyShare => "FOR KEY SHARE",
}
}
fn locking_from_names(stmt: &SelectStatement) -> alloc::vec::Vec<String> {
let mut out = alloc::vec::Vec::new();
if let Some(f) = &stmt.from {
let mut push = |t: &spg_sql::ast::TableRef| {
if let Some(a) = &t.alias {
out.push(a.clone());
}
out.push(t.name.clone());
};
push(&f.primary);
for j in &f.joins {
push(&j.table);
}
}
out
}
fn validate_aggregate_placement(stmt: &SelectStatement) -> Result<(), EngineError> {
use spg_sql::ast::Expr;
if let Some(w) = &stmt.where_
&& aggregate::contains_aggregate(w)
{
return Err(EngineError::Unsupported(
"aggregate functions are not allowed in WHERE".into(),
));
}
let mut nested = false;
let mut check = |e: &Expr| {
let mut probe = e.clone();
crate::expr_analysis::rewrite_nodes_mut(&mut probe, &mut |n| {
let args = match n {
Expr::FunctionCall { name, args } if aggregate::is_aggregate_name(name) => args,
_ => return false,
};
if args.iter().any(aggregate::contains_aggregate) {
nested = true;
}
false
});
};
for it in &stmt.items {
if let spg_sql::ast::SelectItem::Expr { expr, .. } = it {
check(expr);
}
}
if let Some(h) = &stmt.having {
check(h);
}
for o in &stmt.order_by {
check(&o.expr);
}
if nested {
return Err(EngineError::Unsupported(
"aggregate function calls cannot be nested".into(),
));
}
Ok(())
}
fn resolve_positional_order_by(
order_by: &[spg_sql::ast::OrderBy],
projection: &[ProjectedItem],
) -> alloc::vec::Vec<spg_sql::ast::OrderBy> {
order_by
.iter()
.map(|o| {
let mut o = o.clone();
if let Expr::Literal(spg_sql::ast::Literal::Integer(n)) = &o.expr
&& *n >= 1
&& let Ok(idx) = usize::try_from(*n - 1)
&& let Some(item) = projection.get(idx)
&& !expr_contains_builtin_srf(&item.expr)
{
o.expr = item.expr.clone();
}
o
})
.collect()
}
pub(crate) fn expr_contains_builtin_srf(e: &spg_sql::ast::Expr) -> bool {
let mut found = false;
let mut probe = e.clone();
crate::expr_analysis::rewrite_nodes_mut(&mut probe, &mut |n| {
if is_top_level_unnest(n) {
found = true;
return true;
}
false
});
found
}
struct SrfPlan {
nodes: alloc::vec::Vec<spg_sql::ast::Expr>,
rewritten: alloc::vec::Vec<Option<spg_sql::ast::Expr>>,
ext_cols: alloc::vec::Vec<ColumnSchema>,
compiled: alloc::vec::Vec<Option<eval::CompiledExpr>>,
base_cols: usize,
}
fn build_srf_plan(
engine: &Engine,
projection: &[ProjectedItem],
srf_idxs: &[usize],
ctx: &EvalContext<'_>,
) -> Result<SrfPlan, EngineError> {
let mut nodes: Vec<spg_sql::ast::Expr> = Vec::new();
let mut rewritten: Vec<Option<spg_sql::ast::Expr>> = alloc::vec![None; projection.len()];
let mut reject: Option<EngineError> = None;
for &i in srf_idxs {
let mut e = projection[i].expr.clone();
crate::expr_analysis::rewrite_nodes_mut(&mut e, &mut |n| {
if reject.is_some() {
return true;
}
let conditional = match n {
spg_sql::ast::Expr::Case { .. } => Some("CASE"),
spg_sql::ast::Expr::FunctionCall { name, .. }
if name.eq_ignore_ascii_case("coalesce") =>
{
Some("COALESCE")
}
_ => None,
};
if let Some(kind) = conditional
&& engine.expr_contains_srf(n)
{
reject = Some(EngineError::Unsupported(alloc::format!(
"set-returning functions are not allowed in {kind}"
)));
return true;
}
if !engine.is_srf_node(n) {
return false;
}
let slot = nodes.len();
nodes.push(n.clone());
*n = spg_sql::ast::Expr::Column(spg_sql::ast::ColumnName {
qualifier: None,
name: alloc::format!("__srf_{slot}"),
});
true
});
rewritten[i] = Some(e);
}
if let Some(err) = reject {
return Err(err);
}
let base_cols = ctx.columns.len();
let mut ext_cols: Vec<ColumnSchema> = ctx.columns.to_vec();
for slot in 0..nodes.len() {
ext_cols.push(ColumnSchema::new(
alloc::format!("__srf_{slot}"),
DataType::Text,
true,
));
}
let compiled: Vec<Option<eval::CompiledExpr>> = {
let mut ext_ctx = ctx.clone();
ext_ctx.columns = &ext_cols;
projection
.iter()
.enumerate()
.map(|(i, p)| {
let e = rewritten[i].as_ref().unwrap_or(&p.expr);
if eval::fully_compilable(e) {
Some(eval::compile_expr(e, &ext_ctx))
} else {
None
}
})
.collect()
};
Ok(SrfPlan {
nodes,
rewritten,
ext_cols,
compiled,
base_cols,
})
}
fn expand_projection_srfs(
engine: &Engine,
projection: &[ProjectedItem],
srf_idxs: &[usize],
filtered: &[Row<'static>],
ctx: &EvalContext<'_>,
) -> Result<(alloc::vec::Vec<Row<'static>>, alloc::vec::Vec<usize>), EngineError> {
let mut out = alloc::vec::Vec::with_capacity(filtered.len());
let mut src = alloc::vec::Vec::with_capacity(filtered.len());
let mut plan = build_srf_plan(engine, projection, srf_idxs, ctx)?;
let all_pure = projection
.iter()
.enumerate()
.all(|(i, p)| eval::fully_compilable(plan.rewritten[i].as_ref().unwrap_or(&p.expr)))
&& plan.nodes.iter().all(|n| match n {
Expr::FunctionCall { args, .. } => args.iter().all(eval::fully_compilable),
other => eval::fully_compilable(other),
});
if all_pure
&& filtered.len() >= crate::PARALLEL_MIN_ROWS / 5
&& let Some(r) = engine.parallel_runner.0.as_deref()
{
let n_shards = (filtered.len() / (crate::PARALLEL_MIN_ROWS / 5)).clamp(2, 8);
let chunk = filtered.len().div_ceil(n_shards);
type ShardOut = Result<(Vec<Row<'static>>, Vec<usize>), EngineError>;
let schema_cols = ctx.columns;
let alias = ctx.table_alias;
let mysql = ctx.mysql_dialect;
let style = ctx.render_style;
let plan_ref = &plan;
let results = r.run_shards(n_shards, &|si| {
let lo = si * chunk;
let hi = ((si + 1) * chunk).min(filtered.len());
let mut sctx = eval::EvalContext::new(schema_cols, alias);
sctx.mysql_dialect = mysql;
sctx.render_style = style;
let mut local_plan = match build_srf_plan(engine, projection, srf_idxs, &sctx) {
Ok(p) => p,
Err(e) => return alloc::boxed::Box::new(ShardOut::Err(e)) as _,
};
let mut run = || -> ShardOut {
let mut o: Vec<Row<'static>> = Vec::with_capacity(hi - lo);
let mut sidx: Vec<usize> = Vec::with_capacity(hi - lo);
for (i, row) in filtered[lo..hi].iter().enumerate() {
let expanded =
expand_srf_row_with(engine, &mut local_plan, projection, row, &sctx)?;
sidx.extend(core::iter::repeat_n(lo + i, expanded.len()));
o.extend(expanded);
}
Ok((o, sidx))
};
alloc::boxed::Box::new(run())
});
for boxed in results {
let shard = boxed
.downcast::<ShardOut>()
.expect("runner echoes the closure's box");
let (o, sidx) = (*shard)?;
out.extend(o);
src.extend(sidx);
}
return Ok((out, src));
}
for (i, row) in filtered.iter().enumerate() {
let expanded = expand_srf_row_with(engine, &mut plan, projection, row, ctx)?;
src.extend(core::iter::repeat_n(i, expanded.len()));
out.extend(expanded);
}
Ok((out, src))
}
fn srf_order_key(
ob: &spg_sql::ast::OrderBy,
out_col: Option<usize>,
out: &Row<'static>,
src: &Row<'static>,
ctx: &EvalContext<'_>,
) -> Result<Value<'static>, EngineError> {
match out_col {
Some(i) => Ok(out.values.get(i).cloned().unwrap_or(Value::Null)),
None => eval::eval_expr(&ob.expr, src, ctx).map_err(EngineError::Eval),
}
}
fn expand_srf_row_with(
engine: &Engine,
plan: &mut SrfPlan,
projection: &[ProjectedItem],
row: &Row<'static>,
ctx: &EvalContext<'_>,
) -> Result<Vec<Row<'static>>, EngineError> {
let mut lists: Vec<Vec<Value<'static>>> = Vec::with_capacity(plan.nodes.len());
for n in &plan.nodes {
lists.push(engine.srf_values(n, row, ctx)?);
}
let n_rows = lists.iter().map(Vec::len).max().unwrap_or(0);
for (slot, list) in lists.iter().enumerate() {
plan.ext_cols[plan.base_cols + slot].ty = list
.iter()
.find_map(|v| v.data_type())
.unwrap_or(DataType::Text);
}
let mut ext_ctx = ctx.clone();
ext_ctx.columns = &plan.ext_cols;
let mut out = Vec::with_capacity(n_rows);
let base_len = row.values.len();
let mut ext_vals = row.values.clone();
ext_vals.resize(base_len + lists.len(), Value::Null);
let mut eval_stack: alloc::vec::Vec<Value<'static>> = alloc::vec::Vec::new();
for k in 0..n_rows {
for (slot, list) in lists.iter().enumerate() {
ext_vals[base_len + slot] = list.get(k).cloned().unwrap_or(Value::Null);
}
let ext_row = Row::new(core::mem::take(&mut ext_vals));
let mut vals = Vec::with_capacity(projection.len());
for (i, p) in projection.iter().enumerate() {
vals.push(match &plan.compiled[i] {
Some(c) => eval::eval_compiled(c, &ext_row, &ext_ctx, &mut eval_stack)
.map_err(EngineError::Eval)?,
None => {
let expr = plan.rewritten[i].as_ref().unwrap_or(&p.expr);
eval::eval_expr(expr, &ext_row, &ext_ctx).map_err(EngineError::Eval)?
}
});
}
ext_vals = ext_row.values;
out.push(Row::new(vals));
}
Ok(out)
}
fn srf_order_output_cols(
order_by: &[spg_sql::ast::OrderBy],
projection: &[ProjectedItem],
) -> Vec<Option<usize>> {
order_by
.iter()
.map(|ob| {
if let Expr::Literal(spg_sql::ast::Literal::Integer(n)) = &ob.expr
&& *n >= 1
&& let Ok(idx) = usize::try_from(*n - 1)
&& idx < projection.len()
{
return Some(idx);
}
if let Expr::Column(c) = &ob.expr
&& c.qualifier.is_none()
{
let mut hit = None;
for (i, p) in projection.iter().enumerate() {
if p.output_name.eq_ignore_ascii_case(&c.name) {
if hit.is_some() {
hit = None;
break;
}
hit = Some(i);
}
}
if hit.is_some() {
return hit;
}
}
projection.iter().position(|p| p.expr == ob.expr)
})
.collect()
}
fn expand_srf_row(
engine: &Engine,
projection: &[ProjectedItem],
srf_idxs: &[usize],
row: &Row<'static>,
ctx: &EvalContext<'_>,
) -> Result<Vec<Row<'static>>, EngineError> {
let mut plan = build_srf_plan(engine, projection, srf_idxs, ctx)?;
expand_srf_row_with(engine, &mut plan, projection, row, ctx)
}
impl Engine {
fn srf_values(
&self,
expr: &spg_sql::ast::Expr,
row: &Row<'static>,
ctx: &EvalContext<'_>,
) -> Result<Vec<Value<'static>>, EngineError> {
if top_level_srf_kind(expr).is_some() {
return top_level_srf_output(expr, row, ctx);
}
let spg_sql::ast::Expr::FunctionCall { name, args } = expr else {
return Err(EngineError::Unsupported(
"expected a SELECT-list SRF call".into(),
));
};
let mut vals: alloc::vec::Vec<Value<'static>> = alloc::vec::Vec::new();
for a in args {
vals.push(eval::eval_expr(a, row, ctx).map_err(EngineError::Eval)?);
}
let (rows, cols) = self.setof_rows_of(name, &vals, None)?;
Ok(rows
.into_iter()
.map(|r| {
if r.values.len() == 1 {
r.values.into_iter().next().unwrap_or(Value::Null)
} else {
Value::Composite(
cols.iter()
.map(|c| c.name.clone())
.zip(r.values)
.collect::<alloc::vec::Vec<_>>(),
)
}
})
.collect())
}
fn is_srf_node(&self, e: &spg_sql::ast::Expr) -> bool {
if is_top_level_unnest(e) {
return true;
}
let spg_sql::ast::Expr::FunctionCall { name, .. } = e else {
return false;
};
self.active_catalog().functions_named(name).iter().any(|f| {
let r = f.returns.trim().to_ascii_uppercase();
r.starts_with("SETOF") || r.starts_with("TABLE(")
})
}
fn expr_contains_srf(&self, e: &spg_sql::ast::Expr) -> bool {
let mut found = false;
let mut probe = e.clone();
crate::expr_analysis::rewrite_nodes_mut(&mut probe, &mut |n| {
if self.is_srf_node(n) {
found = true;
return true;
}
false
});
found
}
fn srf_target_idxs(&self, projection: &[ProjectedItem]) -> alloc::vec::Vec<usize> {
projection
.iter()
.enumerate()
.filter(|(_, p)| self.expr_contains_srf(&p.expr))
.map(|(i, _)| i)
.collect()
}
}
impl Engine {
fn lower_record_expansion(
&self,
stmt: &SelectStatement,
) -> Result<Option<SelectStatement>, EngineError> {
use spg_sql::ast::{Expr, SelectItem};
let is_marker = |it: &SelectItem| {
matches!(it, SelectItem::Expr { expr: Expr::FunctionCall { name, .. }, .. }
if name == "__record_expand")
};
if !stmt.items.iter().any(is_marker) {
return Ok(None);
}
let mut out = stmt.clone();
let mut items: alloc::vec::Vec<SelectItem> = alloc::vec::Vec::new();
let mut lateral_refs: alloc::vec::Vec<TableRef> = alloc::vec::Vec::new();
for (n, item) in stmt.items.iter().enumerate() {
if !is_marker(item) {
items.push(item.clone());
continue;
}
let SelectItem::Expr {
expr: Expr::FunctionCall { args, .. },
..
} = item
else {
unreachable!("checked by is_marker");
};
let Some(Expr::FunctionCall {
name: fname,
args: fargs,
}) = args.first()
else {
return Err(EngineError::Unsupported(
"(<expr>).* expands a function's record — it needs a function call".into(),
));
};
let cols = self.setof_declared_columns(fname)?;
let alias = alloc::format!("__rec{n}");
let mut tref = bare_table_ref_named(&alias);
tref.table_fn_call = Some(alloc::boxed::Box::new((
fname.to_ascii_lowercase(),
fargs.clone(),
)));
tref.alias = Some(alias.clone());
lateral_refs.push(tref);
for c in cols {
items.push(SelectItem::Expr {
expr: Expr::Column(spg_sql::ast::ColumnName {
qualifier: Some(alias.clone()),
name: c,
}),
alias: None,
});
}
}
out.items = items;
for tref in lateral_refs {
match &mut out.from {
None => {
out.from = Some(spg_sql::ast::FromClause {
primary: tref,
joins: alloc::vec::Vec::new(),
});
}
Some(from) => from.joins.push(spg_sql::ast::FromJoin {
kind: spg_sql::ast::JoinKind::Cross,
table: tref,
on: None,
using_cols: None,
natural: false,
}),
}
}
Ok(Some(out))
}
fn setof_declared_columns(
&self,
name: &str,
) -> Result<alloc::vec::Vec<alloc::string::String>, EngineError> {
let cat = self.active_catalog();
let overloads = cat.functions_named(name);
let def = overloads.first().ok_or_else(|| {
EngineError::Unsupported(alloc::format!("function {name} does not exist"))
})?;
let declared = def.returns.trim();
let upper = declared.to_ascii_uppercase();
if upper.starts_with("TABLE(") {
let raw = &declared["TABLE(".len()..declared.len() - 1];
return Ok(raw
.split(',')
.map(|d| d.split_whitespace().next().unwrap_or("col").to_string())
.collect());
}
Ok(alloc::vec![name.to_string()])
}
}
pub(crate) fn json_table_schema_pub(
cols: &[spg_sql::ast::JsonTableColumn],
) -> alloc::vec::Vec<ColumnSchema> {
json_table_schema(cols)
}
fn json_table_schema(cols: &[spg_sql::ast::JsonTableColumn]) -> alloc::vec::Vec<ColumnSchema> {
use spg_sql::ast::JsonTableColumn as C;
let mut out = alloc::vec::Vec::new();
for c in cols {
match c {
C::Ordinality { name } => {
out.push(ColumnSchema::new(name.clone(), DataType::BigInt, false));
}
C::Regular {
name, ty, exists, ..
} => {
let dt = if *exists {
DataType::Bool
} else {
crate::conversions::column_type_to_data_type(*ty)
};
out.push(ColumnSchema::new(name.clone(), dt, true));
}
C::Nested { columns, .. } => out.extend(json_table_schema(columns)),
}
}
out
}
fn coerce_json_table_default(
v: Value<'static>,
ty: spg_sql::ast::ColumnTypeName,
name: &str,
) -> Result<Value<'static>, EngineError> {
if v.is_null() {
return Ok(Value::Null);
}
let dt = crate::conversions::column_type_to_data_type(ty);
crate::conversions::coerce_value(v, dt, name, 0)
}
fn value_to_json_value(v: &Value<'_>) -> crate::json::JsonValue {
use crate::json::JsonValue as J;
match v {
Value::Null => J::Null,
Value::Bool(b) => J::Bool(*b),
Value::SmallInt(n) => J::Number(f64::from(*n)),
Value::Int(n) => J::Number(f64::from(*n)),
Value::BigInt(n) => J::Number(*n as f64),
Value::Float(x) => J::Number(*x),
Value::Json(s) => crate::json::parse_doc(s).unwrap_or(J::Null),
other => J::String(crate::eval::value_to_text(other)),
}
}
fn bare_table_ref_named(name: &str) -> TableRef {
TableRef {
name: name.to_string(),
alias: None,
only: false,
as_of_segment: None,
unnest_expr: None,
unnest_column_aliases: alloc::vec::Vec::new(),
with_ordinality: false,
generate_series_args: None,
lateral_subquery: None,
jsonb_each_text_arg: None,
table_fn_call: None,
rows_from: None,
json_table: None,
scalar_fn_item: false,
}
}
impl Engine {
fn rows_from_rows(
&self,
primary: &TableRef,
) -> Result<(alloc::vec::Vec<Row<'static>>, alloc::vec::Vec<ColumnSchema>), EngineError> {
let entries = primary
.rows_from
.as_ref()
.expect("caller guards rows_from.is_some()");
let empty: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
let ctx = self.ev_ctx(&empty, None);
let dummy = Row::new(alloc::vec::Vec::new());
let mut lists: alloc::vec::Vec<alloc::vec::Vec<Value<'static>>> = alloc::vec::Vec::new();
let mut cols: alloc::vec::Vec<ColumnSchema> = alloc::vec::Vec::new();
for (name, args) in entries {
let (vals, colname) = if name == "__array" {
let arr = eval::eval_expr(&args[0], &dummy, &ctx).map_err(EngineError::Eval)?;
(
array_value_to_elements(&arr)?,
alloc::string::String::from("unnest"),
)
} else {
let call = spg_sql::ast::Expr::FunctionCall {
name: name.clone(),
args: args.clone(),
};
(self.srf_values(&call, &dummy, &ctx)?, name.clone())
};
let ty = vals
.first()
.and_then(spg_storage::Value::data_type)
.unwrap_or(DataType::Text);
cols.push(ColumnSchema::new(colname, ty, true));
lists.push(vals);
}
let n = lists.iter().map(alloc::vec::Vec::len).max().unwrap_or(0);
let mut rows: alloc::vec::Vec<Row<'static>> = alloc::vec::Vec::with_capacity(n);
for k in 0..n {
let mut vals: alloc::vec::Vec<Value<'static>> =
alloc::vec::Vec::with_capacity(lists.len() + 1);
for l in &lists {
vals.push(l.get(k).cloned().unwrap_or(Value::Null));
}
rows.push(Row::new(vals));
}
if primary.with_ordinality {
cols.push(ColumnSchema::new(
"ordinality".to_string(),
DataType::BigInt,
false,
));
rows = rows
.into_iter()
.enumerate()
.map(|(i, r)| {
let mut v = r.values;
v.push(Value::BigInt(i as i64 + 1));
Row::new(v)
})
.collect();
}
Ok((rows, cols))
}
}
fn set_op_name(kind: UnionKind) -> &'static str {
match kind {
UnionKind::All | UnionKind::Distinct => "UNION",
UnionKind::Intersect | UnionKind::IntersectAll => "INTERSECT",
UnionKind::Except | UnionKind::ExceptAll => "EXCEPT",
}
}
fn branch_unknown_mask(stmt: &SelectStatement) -> Vec<bool> {
stmt.items
.iter()
.map(|item| match item {
SelectItem::Expr { expr, .. } => matches!(
expr,
Expr::Literal(spg_sql::ast::Literal::String(_))
| Expr::Literal(spg_sql::ast::Literal::Null)
),
_ => false,
})
.collect()
}
fn coerce_branch_column(
rows: &mut [Row<'static>],
col_idx: usize,
target: DataType,
col_name: &str,
) -> Result<(), EngineError> {
for row in rows.iter_mut() {
let Some(slot) = row.values.get_mut(col_idx) else {
continue;
};
if matches!(slot, Value::Null) {
continue;
}
*slot = crate::conversions::coerce_value(slot.clone(), target, col_name, col_idx)?;
}
Ok(())
}
fn try_flatten_derived(stmt: &SelectStatement, primary: &TableRef) -> Option<SelectStatement> {
use spg_sql::ast::SelectItem;
let inner = primary.lateral_subquery.as_deref()?;
if !stmt.ctes.is_empty()
|| !stmt.unions.is_empty()
|| stmt.distinct
|| !stmt.distinct_on.is_empty()
|| !stmt.window_check_exprs.is_empty()
|| stmt.locking.is_some()
|| primary.with_ordinality
|| !primary.unnest_column_aliases.is_empty()
{
return None;
}
if !inner.ctes.is_empty()
|| !inner.unions.is_empty()
|| inner.distinct
|| !inner.distinct_on.is_empty()
|| inner.group_by.is_some()
|| inner.group_by_all
|| inner.having.is_some()
|| !inner.order_by.is_empty()
|| inner.limit.is_some()
|| inner.offset.is_some()
|| !inner.window_check_exprs.is_empty()
|| inner.locking.is_some()
{
return None;
}
let ifrom = inner.from.as_ref()?;
let it = &ifrom.primary;
if !ifrom.joins.is_empty()
|| it.name.is_empty()
|| it.lateral_subquery.is_some()
|| it.unnest_expr.is_some()
|| it.generate_series_args.is_some()
|| it.as_of_segment.is_some()
|| it.jsonb_each_text_arg.is_some()
|| it.table_fn_call.is_some()
|| it.rows_from.is_some()
|| it.json_table.is_some()
|| it.with_ordinality
|| !it.unnest_column_aliases.is_empty()
{
return None;
}
if inner.where_.as_ref().is_some_and(crate::expr_has_subquery) {
return None;
}
let inner_alias = it.alias.clone().unwrap_or_else(|| it.name.clone());
let mut map: alloc::collections::BTreeMap<String, spg_sql::ast::ColumnName> =
alloc::collections::BTreeMap::new();
for item in &inner.items {
let SelectItem::Expr { expr, alias } = item else {
return None;
};
let Expr::Column(c) = expr else {
return None;
};
if let Some(q) = c.qualifier.as_deref()
&& !q.eq_ignore_ascii_case(&inner_alias)
{
return None;
}
let out_name = alias.clone().unwrap_or_else(|| c.name.clone());
if map
.insert(out_name.to_ascii_lowercase(), c.clone())
.is_some()
{
return None;
}
}
if map.is_empty() {
return None;
}
let derived_alias = primary
.alias
.clone()
.unwrap_or_else(|| primary.name.clone())
.to_ascii_lowercase();
let mut out = stmt.clone();
let ok = core::cell::Cell::new(true);
let mut subst = |e: &mut Expr| -> bool {
match e {
Expr::Column(c) => {
match c.qualifier.as_deref() {
Some(q) if q.eq_ignore_ascii_case(&derived_alias) => {}
None => {}
Some(_) => {
ok.set(false);
return true;
}
}
match map.get(&c.name.to_ascii_lowercase()) {
Some(target) => *c = target.clone(),
None => ok.set(false),
}
true
}
Expr::ScalarSubquery(_)
| Expr::Exists { .. }
| Expr::InSubquery { .. }
| Expr::RowInSubquery { .. }
| Expr::RowCmpSubquery { .. } => {
ok.set(false);
true
}
_ => false,
}
};
for item in &mut out.items {
match item {
SelectItem::Expr { expr, .. } => {
crate::expr_analysis::rewrite_nodes_mut(expr, &mut subst);
}
SelectItem::Wildcard | SelectItem::QualifiedWildcard(_) => return None,
}
}
if let Some(w) = &mut out.where_ {
crate::expr_analysis::rewrite_nodes_mut(w, &mut subst);
}
if let Some(gs) = &mut out.group_by {
for g in gs {
crate::expr_analysis::rewrite_nodes_mut(g, &mut subst);
}
}
if let Some(h) = &mut out.having {
crate::expr_analysis::rewrite_nodes_mut(h, &mut subst);
}
for o in &mut out.order_by {
crate::expr_analysis::rewrite_nodes_mut(&mut o.expr, &mut subst);
}
for d in &mut out.distinct_on {
crate::expr_analysis::rewrite_nodes_mut(d, &mut subst);
}
if !ok.get() {
return None;
}
out.from = Some(spg_sql::ast::FromClause {
primary: it.clone(),
joins: Vec::new(),
});
out.where_ = match (inner.where_.clone(), out.where_.take()) {
(Some(a), Some(b)) => Some(Expr::Binary {
lhs: alloc::boxed::Box::new(a),
op: spg_sql::ast::BinOp::And,
rhs: alloc::boxed::Box::new(b),
}),
(Some(a), None) => Some(a),
(None, b) => b,
};
Some(out)
}
fn try_count_over_offset(stmt: &SelectStatement, primary: &TableRef) -> Option<SelectStatement> {
use spg_sql::ast::{Expr as E, LimitExpr, SelectItem};
let inner = primary.lateral_subquery.as_deref()?;
if !stmt.ctes.is_empty()
|| !stmt.unions.is_empty()
|| stmt.distinct
|| !stmt.distinct_on.is_empty()
|| stmt.where_.is_some()
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| !stmt.order_by.is_empty()
|| stmt.limit.is_some()
|| stmt.offset.is_some()
|| stmt.items.len() != 1
{
return None;
}
let SelectItem::Expr { expr, .. } = &stmt.items[0] else {
return None;
};
let E::FunctionCall { name, args } = expr else {
return None;
};
if !name.eq_ignore_ascii_case("count_star") || !args.is_empty() {
return None;
}
let Some(LimitExpr::Literal(k)) = &inner.offset else {
return None;
};
let k = i64::from(*k);
if inner.limit.is_some() || inner.order_by.is_empty() {
return None;
}
let mut counted = inner.clone();
counted.order_by = Vec::new();
counted.offset = None;
let base = matview_flatten_probe(&counted)?;
let mut out = stmt.clone();
out.items = alloc::vec![SelectItem::Expr {
expr: E::FunctionCall {
name: String::from("greatest"),
args: alloc::vec![
E::Binary {
lhs: alloc::boxed::Box::new(E::FunctionCall {
name: String::from("count_star"),
args: alloc::vec![],
}),
op: spg_sql::ast::BinOp::Sub,
rhs: alloc::boxed::Box::new(E::Literal(spg_sql::ast::Literal::Integer(k))),
},
E::Literal(spg_sql::ast::Literal::Integer(0)),
],
},
alias: Some(String::from("count")),
}];
out.from = Some(spg_sql::ast::FromClause {
primary: base,
joins: Vec::new(),
});
out.where_ = counted.where_.clone();
Some(out)
}
fn matview_flatten_probe(inner: &SelectStatement) -> Option<TableRef> {
use spg_sql::ast::SelectItem;
if !inner.ctes.is_empty()
|| !inner.unions.is_empty()
|| inner.distinct
|| !inner.distinct_on.is_empty()
|| inner.group_by.is_some()
|| inner.group_by_all
|| inner.having.is_some()
|| !inner.order_by.is_empty()
|| inner.limit.is_some()
|| inner.offset.is_some()
|| !inner.window_check_exprs.is_empty()
|| inner.locking.is_some()
{
return None;
}
let ifrom = inner.from.as_ref()?;
let it = &ifrom.primary;
if !ifrom.joins.is_empty()
|| it.name.is_empty()
|| it.lateral_subquery.is_some()
|| it.unnest_expr.is_some()
|| it.generate_series_args.is_some()
|| it.as_of_segment.is_some()
|| it.jsonb_each_text_arg.is_some()
|| it.table_fn_call.is_some()
|| it.rows_from.is_some()
|| it.json_table.is_some()
|| it.with_ordinality
{
return None;
}
for item in &inner.items {
match item {
SelectItem::Expr { expr, .. } => {
if crate::expr_has_subquery(expr) || expr_contains_builtin_srf(expr) {
return None;
}
}
SelectItem::Wildcard => {}
SelectItem::QualifiedWildcard(_) => return None,
}
}
if inner.where_.as_ref().is_some_and(crate::expr_has_subquery) {
return None;
}
Some(it.clone())
}
fn try_count_over_const_unnest(
stmt: &SelectStatement,
primary: &TableRef,
) -> Option<SelectStatement> {
use spg_sql::ast::{Expr as E, SelectItem};
let inner = primary.lateral_subquery.as_deref()?;
if !stmt.ctes.is_empty()
|| !stmt.unions.is_empty()
|| stmt.distinct
|| !stmt.distinct_on.is_empty()
|| stmt.where_.is_some()
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| !stmt.order_by.is_empty()
|| stmt.limit.is_some()
|| stmt.offset.is_some()
|| stmt.items.len() != 1
{
return None;
}
let SelectItem::Expr { expr, .. } = &stmt.items[0] else {
return None;
};
let E::FunctionCall { name, args } = expr else {
return None;
};
if !name.eq_ignore_ascii_case("count_star") || !args.is_empty() {
return None;
}
if inner.items.len() != 1
|| !inner.order_by.is_empty()
|| inner.limit.is_some()
|| inner.offset.is_some()
{
return None;
}
let SelectItem::Expr { expr: item, .. } = &inner.items[0] else {
return None;
};
let E::FunctionCall {
name: fname,
args: fargs,
} = item
else {
return None;
};
if !fname.eq_ignore_ascii_case("unnest") || fargs.len() != 1 {
return None;
}
let E::Array(elems) = &fargs[0] else {
return None;
};
if elems.is_empty() || elems.iter().any(crate::expr_has_subquery) {
return None;
}
let k = elems.len() as i64;
let mut counted = inner.clone();
counted.items = alloc::vec![SelectItem::Expr {
expr: E::Literal(spg_sql::ast::Literal::Integer(1)),
alias: None,
}];
let base = matview_flatten_probe(&counted)?;
let mut out = stmt.clone();
out.items = alloc::vec![SelectItem::Expr {
expr: E::Binary {
lhs: alloc::boxed::Box::new(E::FunctionCall {
name: String::from("count_star"),
args: alloc::vec![],
}),
op: spg_sql::ast::BinOp::Mul,
rhs: alloc::boxed::Box::new(E::Literal(spg_sql::ast::Literal::Integer(k))),
},
alias: Some(String::from("count")),
}];
out.from = Some(spg_sql::ast::FromClause {
primary: base,
joins: Vec::new(),
});
out.where_ = counted.where_.clone();
Some(out)
}