use alloc::borrow::Cow;
use alloc::string::{String, ToString};
use alloc::vec::Vec;
use spg_sql::ast::{Expr, FromClause, JoinKind, SelectItem, SelectStatement, TableRef};
use spg_storage::{ColumnSchema, DataType, Row, Table, Value};
use crate::eval::EvalContext;
use crate::{
ByteBudget, CancelToken, Engine, EngineError, OrderKey, QueryResult, aggregate,
apply_offset_and_limit, approx_row_bytes, approx_rows_bytes, approx_value_bytes,
build_order_keys, build_projection, cmp_multi_key, collect_column_qualifiers,
collect_qualified_refs, eval, expr_has_subquery, memoize, reorder, value_cmp,
value_to_literal_expr,
};
pub(crate) struct JoinedPeer<'a> {
pub(crate) eager_rows: Option<Vec<Row<'static>>>,
pub(crate) cols: Vec<ColumnSchema>,
pub(crate) alias: String,
pub(crate) kind: JoinKind,
pub(crate) on: Option<&'a Expr>,
pub(crate) lateral: Option<&'a SelectStatement>,
pub(crate) join_table: Option<String>,
pub(crate) where_preds: Vec<Expr>,
}
pub(crate) enum JoinSrc<'a> {
Owned(Vec<Row<'static>>),
Eager(&'a [Row<'static>]),
Stored(&'a spg_storage::persistent::PersistentVec<Row<'static>>),
Mixed {
hot: &'a spg_storage::persistent::PersistentVec<Row<'static>>,
cold: Vec<Row<'static>>,
cold_locator_map: hashbrown::HashMap<i64, usize>,
},
}
#[derive(Debug)]
enum Bucket {
One(usize),
Many(Vec<usize>),
}
impl Bucket {
fn push(&mut self, ri: usize) {
match self {
Self::One(first) => *self = Self::Many(alloc::vec![*first, ri]),
Self::Many(v) => v.push(ri),
}
}
fn as_slice(&self) -> &[usize] {
match self {
Self::One(x) => core::slice::from_ref(x),
Self::Many(v) => v.as_slice(),
}
}
}
impl JoinSrc<'_> {
pub(crate) fn get(&self, i: usize) -> Option<&Row<'static>> {
match self {
Self::Owned(v) => v.get(i),
Self::Eager(s) => s.get(i),
Self::Stored(p) => p.get(i),
Self::Mixed { hot, cold, .. } => {
if i < hot.len() {
hot.get(i)
} else {
cold.get(i - hot.len())
}
}
}
}
pub(crate) fn len(&self) -> usize {
match self {
Self::Owned(v) => v.len(),
Self::Eager(s) => s.len(),
Self::Stored(p) => p.len(),
Self::Mixed { hot, cold, .. } => hot.len() + cold.len(),
}
}
pub(crate) fn cold_pk_offset(&self, pk_key: i64) -> Option<usize> {
match self {
Self::Mixed {
hot,
cold_locator_map,
..
} => cold_locator_map
.get(&pk_key)
.copied()
.map(|off| hot.len() + off),
_ => None,
}
}
}
pub(crate) fn tuple_value<'s>(
sources: &'s [JoinSrc<'_>],
offsets: &[usize],
tuple: &[usize],
pos: usize,
) -> Option<&'s Value<'static>> {
let k = offsets.partition_point(|&o| o <= pos).checked_sub(1)?;
let ri = *tuple.get(k)?;
if ri == usize::MAX {
return None;
}
sources.get(k)?.get(ri)?.values.get(pos - offsets[k])
}
#[inline]
pub(crate) fn tuple_value_indexed<'s>(
sources: &'s [JoinSrc<'_>],
offsets: &[usize],
pos_to_src: &[u16],
tuple: &[usize],
pos: usize,
) -> Option<&'s Value<'static>> {
let k = *pos_to_src.get(pos)? as usize;
let ri = *tuple.get(k)?;
if ri == usize::MAX {
return None;
}
sources.get(k)?.get(ri)?.values.get(pos - offsets[k])
}
pub(crate) fn build_pos_to_src(offsets: &[usize]) -> Vec<u16> {
let width = offsets.last().copied().unwrap_or(0);
let mut tab: Vec<u16> = Vec::with_capacity(width);
for k in 0..offsets.len().saturating_sub(1) {
let span = offsets[k + 1] - offsets[k];
for _ in 0..span {
tab.push(k as u16);
}
}
tab
}
#[derive(Clone, Copy)]
pub(crate) enum AggRows<'a> {
Owned(&'a [Row<'static>]),
Refs(&'a [RowRef<'a>]),
Ptrs(&'a [&'a Row<'static>]),
}
impl<'a> AggRows<'a> {
#[inline]
pub(crate) fn len(&self) -> usize {
match self {
Self::Owned(r) => r.len(),
Self::Refs(r) => r.len(),
Self::Ptrs(r) => r.len(),
}
}
#[inline]
pub(crate) fn is_empty(&self) -> bool {
self.len() == 0
}
#[inline]
pub(crate) fn get(&self, i: usize) -> Option<RowRef<'a>> {
match self {
Self::Owned(r) => r.get(i).map(RowRef::Owned),
Self::Refs(r) => r.get(i).copied(),
Self::Ptrs(r) => r.get(i).map(|p| RowRef::Owned(p)),
}
}
#[inline]
pub(crate) fn range(&self, lo: usize, hi: usize) -> Self {
match self {
Self::Owned(r) => Self::Owned(&r[lo..hi]),
Self::Refs(r) => Self::Refs(&r[lo..hi]),
Self::Ptrs(r) => Self::Ptrs(&r[lo..hi]),
}
}
#[inline]
pub(crate) fn first(&self) -> Option<RowRef<'a>> {
self.get(0)
}
#[inline]
pub(crate) fn iter(&self) -> impl Iterator<Item = RowRef<'a>> + '_ {
(0..self.len()).filter_map(move |i| self.get(i))
}
}
#[derive(Clone, Copy)]
pub(crate) enum RowRef<'a> {
Owned(&'a Row<'static>),
Tuple {
sources: &'a [JoinSrc<'a>],
offsets: &'a [usize],
pos_to_src: &'a [u16],
tuple: &'a [usize],
},
}
impl<'a> RowRef<'a> {
#[inline]
pub(crate) fn get(&self, pos: usize) -> Option<&'a Value<'a>> {
match self {
RowRef::Owned(r) => r.values.get(pos),
RowRef::Tuple {
sources,
offsets,
pos_to_src,
tuple,
} => tuple_value_indexed(sources, offsets, pos_to_src, tuple, pos),
}
}
pub(crate) fn as_row(&self) -> Cow<'_, Row<'static>> {
match self {
RowRef::Owned(r) => Cow::Borrowed(r),
RowRef::Tuple {
sources,
offsets,
pos_to_src,
tuple,
} => {
let width = offsets.last().copied().unwrap_or(0);
let mut vals: Vec<Value<'static>> = Vec::with_capacity(width);
for pos in 0..width {
vals.push(
tuple_value_indexed(sources, offsets, pos_to_src, tuple, pos)
.cloned()
.unwrap_or(Value::Null),
);
}
Cow::Owned(Row::new(vals))
}
}
}
pub(crate) fn as_row_into(&self, buf: &mut Vec<Value<'static>>) {
buf.clear();
match self {
RowRef::Owned(r) => {
buf.reserve(r.values.len());
for v in &r.values {
buf.push(v.clone());
}
}
RowRef::Tuple {
sources,
offsets,
pos_to_src,
tuple,
} => {
let width = offsets.last().copied().unwrap_or(0);
buf.reserve(width);
for pos in 0..width {
buf.push(
tuple_value_indexed(sources, offsets, pos_to_src, tuple, pos)
.cloned()
.unwrap_or(Value::Null),
);
}
}
}
}
}
pub(crate) fn extend_masked(
vals: &mut Vec<Value<'static>>,
row: &Row<'static>,
mask: Option<&[bool]>,
) {
match mask {
Some(keep) => {
for (i, v) in row.values.iter().enumerate() {
if keep.get(i).copied().unwrap_or(false) {
vals.push(v.clone());
} else {
vals.push(Value::Null);
}
}
}
None => vals.extend(row.values.iter().cloned()),
}
}
pub(crate) fn materialise_tuple_vals(
sources: &[JoinSrc<'_>],
widths: &[usize],
masks: &[Option<Vec<bool>>],
tuple: &[usize],
cap: usize,
) -> Vec<Value<'static>> {
let mut vals: Vec<Value<'static>> = Vec::with_capacity(cap);
for (k, &ri) in tuple.iter().enumerate() {
let row = if ri == usize::MAX {
None
} else {
sources[k].get(ri)
};
match row {
Some(r) => extend_masked(&mut vals, r, masks[k].as_deref()),
None => {
for _ in 0..widths[k] {
vals.push(Value::Null);
}
}
}
}
vals
}
pub(crate) struct DeferredJoin<'a> {
pub(crate) sources: Vec<JoinSrc<'a>>,
pub(crate) offsets: Vec<usize>,
pub(crate) pos_to_src: Vec<u16>,
pub(crate) widths: Vec<usize>,
pub(crate) masks: Vec<Option<Vec<bool>>>,
pub(crate) survivors: Vec<usize>,
pub(crate) stride: usize,
pub(crate) combined_schema: Vec<ColumnSchema>,
}
impl DeferredJoin<'_> {
pub(crate) fn len(&self) -> usize {
if self.stride == 0 {
0
} else {
self.survivors.len() / self.stride
}
}
pub(crate) fn row_refs(&self) -> Vec<RowRef<'_>> {
if self.stride == 0 {
return Vec::new();
}
self.survivors
.chunks(self.stride)
.map(|tuple| RowRef::Tuple {
sources: &self.sources,
offsets: &self.offsets,
pos_to_src: &self.pos_to_src,
tuple,
})
.collect()
}
pub(crate) fn materialise(&self) -> Vec<Row<'static>> {
if self.stride == 0 {
return Vec::new();
}
let cap = self.offsets.last().copied().unwrap_or(0);
self.survivors
.chunks(self.stride)
.map(|tuple| {
Row::new(materialise_tuple_vals(
&self.sources,
&self.widths,
&self.masks,
tuple,
cap,
))
})
.collect()
}
}
pub(crate) fn approx_tuple_bytes(
sources: &[JoinSrc<'_>],
offsets: &[usize],
masks: &[Option<Vec<bool>>],
tuple: &[usize],
) -> usize {
let width = offsets.last().copied().unwrap_or(0);
let mut bytes = width * core::mem::size_of::<Value>();
for (k, &ri) in tuple.iter().enumerate() {
if ri == usize::MAX {
continue;
}
let Some(row) = sources.get(k).and_then(|s| s.get(ri)) else {
continue;
};
let mask = masks.get(k).and_then(|m| m.as_deref());
for (i, v) in row.values.iter().enumerate() {
let kept = mask.map_or(true, |m| m.get(i).copied().unwrap_or(false));
if kept {
bytes += approx_value_bytes(v);
}
}
}
bytes
}
struct TopNEntry {
keys: Vec<OrderKey>,
descs: alloc::rc::Rc<[bool]>,
seq: u64,
row: Row<'static>,
}
impl PartialEq for TopNEntry {
fn eq(&self, other: &Self) -> bool {
self.cmp(other) == core::cmp::Ordering::Equal
}
}
impl Eq for TopNEntry {}
impl PartialOrd for TopNEntry {
fn partial_cmp(&self, other: &Self) -> Option<core::cmp::Ordering> {
Some(self.cmp(other))
}
}
impl Ord for TopNEntry {
fn cmp(&self, other: &Self) -> core::cmp::Ordering {
cmp_multi_key(&self.keys, &other.keys, &self.descs).then(self.seq.cmp(&other.seq))
}
}
const MAX_JOIN_INTERMEDIATE_ROWS: usize = 4_000_000;
struct JoinPipeline<'a> {
sources: Vec<JoinSrc<'a>>,
masks: Vec<Option<Vec<bool>>>,
widths: Vec<usize>,
offsets: Vec<usize>,
pos_to_src: Vec<u16>,
working: Vec<usize>,
stride: usize,
consumed_cols: usize,
}
impl<'a> JoinPipeline<'a> {
fn new(
primary: JoinSrc<'a>,
mask: Option<Vec<bool>>,
width: usize,
working: Vec<usize>,
) -> Self {
let offsets = alloc::vec![0, width];
let pos_to_src = build_pos_to_src(&offsets);
Self {
sources: alloc::vec![primary],
masks: alloc::vec![mask],
widths: alloc::vec![width],
offsets,
pos_to_src,
working,
stride: 1,
consumed_cols: width,
}
}
fn rows(&self) -> usize {
self.working.len() / self.stride
}
fn advance(
&mut self,
next: Vec<usize>,
source: JoinSrc<'a>,
mask: Option<Vec<bool>>,
right_arity: usize,
) {
self.working = next;
self.stride += 1;
self.sources.push(source);
self.masks.push(mask);
self.consumed_cols += right_arity;
self.offsets.push(self.consumed_cols);
self.widths.push(right_arity);
let k = (self.sources.len() - 1) as u16;
for _ in 0..right_arity {
self.pos_to_src.push(k);
}
}
}
fn keep_mask(
needed: Option<&alloc::collections::BTreeSet<(String, String)>>,
cols: &[ColumnSchema],
alias: &str,
) -> Option<Vec<bool>> {
let needed = needed?;
let keep: Vec<bool> = cols
.iter()
.map(|c| needed.contains(&(alias.to_string(), c.name.clone())))
.collect();
if keep.iter().all(|k| *k) {
None
} else {
Some(keep)
}
}
fn extract_join_keys<'a>(
peer: &JoinedPeer<'a>,
combined_schema: &[ColumnSchema],
consumed_cols: usize,
) -> (
Vec<(usize, usize)>,
// v7.39 (round 719) — the third member is the whole CONJUNCT the
// (left-pos, key-expr) pair came from, so the int-keyed lane can
// identify it in `residual` and drop the re-verification (see
// `join_stage_hash`).
Vec<(usize, &'a Expr, &'a Expr)>,
// v7.39 (round 720) — the MIRROR: `<peer column> = <integer-only
// expression over the joined left side>` (`ON b.id = a.id + 500000`,
// the shape the EXISTS pull-up emits). (peer-col pos, left expr,
// conjunct). Only the integer-only shape is collected — anything
// else keeps the residual path it has today.
Vec<(usize, &'a Expr, &'a Expr)>,
Vec<&'a Expr>,
) {
let mut eq_pairs: Vec<(usize, usize)> = Vec::new();
let mut eq_exprs: Vec<(usize, &Expr, &Expr)> = Vec::new();
let mut eq_probe_exprs: Vec<(usize, &Expr, &Expr)> = Vec::new();
let mut residual: Vec<&Expr> = Vec::new();
if let (Some(on_expr), None) = (peer.on, peer.lateral) {
for sub in reorder::split_and_conjunctions(on_expr) {
if let Some(pair) = match_equi_pair(sub, peer, combined_schema, consumed_cols) {
eq_pairs.push(pair);
continue;
}
if let Some((l, e)) = match_equi_expr(sub, peer, combined_schema, consumed_cols) {
eq_exprs.push((l, e, sub));
residual.push(sub);
continue;
}
if let Some((p, e)) = match_equi_probe_expr(sub, peer, combined_schema, consumed_cols) {
eq_probe_exprs.push((p, e, sub));
residual.push(sub);
continue;
}
residual.push(sub);
}
}
(eq_pairs, eq_exprs, eq_probe_exprs, residual)
}
fn match_equi_probe_expr<'a>(
sub: &'a Expr,
peer: &JoinedPeer<'_>,
combined_schema: &[ColumnSchema],
consumed_cols: usize,
) -> Option<(usize, &'a Expr)> {
let Expr::Binary {
lhs,
op: spg_sql::ast::BinOp::Eq,
rhs,
} = sub
else {
return None;
};
let left_slice = &combined_schema[..consumed_cols];
for (a, b) in [(lhs.as_ref(), rhs.as_ref()), (rhs.as_ref(), lhs.as_ref())] {
if let Expr::Column(c) = a
&& let Some(p) = Engine::peer_col_pos(&peer.alias, &peer.cols, c)
&& matches!(
peer.cols[p].ty,
spg_storage::DataType::Int
| spg_storage::DataType::BigInt
| spg_storage::DataType::SmallInt
)
&& !matches!(b, Expr::Column(_))
&& expr_mentions_a_column(b)
&& int_only_left_expr(b, left_slice)
{
return Some((p, b));
}
}
None
}
fn match_equi_expr<'a>(
sub: &'a Expr,
peer: &JoinedPeer<'_>,
combined_schema: &[ColumnSchema],
consumed_cols: usize,
) -> Option<(usize, &'a Expr)> {
let Expr::Binary {
lhs,
op: spg_sql::ast::BinOp::Eq,
rhs,
} = sub
else {
return None;
};
let left_slice = &combined_schema[..consumed_cols];
for (a, b) in [(lhs.as_ref(), rhs.as_ref()), (rhs.as_ref(), lhs.as_ref())] {
if let Expr::Column(c) = a
&& let Some(l) = Engine::composite_col_pos(left_slice, c)
&& !matches!(b, Expr::Column(_))
&& peer_only_key_expr(b, peer)
&& expr_mentions_a_column(b)
{
return Some((l, b));
}
}
None
}
fn peer_only_key_expr(e: &Expr, peer: &JoinedPeer<'_>) -> bool {
use spg_sql::ast::BinOp;
match e {
Expr::Column(c) => Engine::peer_col_pos(&peer.alias, &peer.cols, c).is_some(),
Expr::Literal(_) => true,
Expr::Unary { expr, .. } | Expr::Cast { expr, .. } => peer_only_key_expr(expr, peer),
Expr::Binary { lhs, op, rhs } => {
matches!(
op,
BinOp::Add | BinOp::Sub | BinOp::Mul | BinOp::Div | BinOp::IntDiv | BinOp::Mod
) && peer_only_key_expr(lhs, peer)
&& peer_only_key_expr(rhs, peer)
}
_ => false,
}
}
fn int_only_key_expr(e: &Expr, peer: &JoinedPeer<'_>) -> bool {
use spg_sql::ast::BinOp;
match e {
Expr::Column(c) => Engine::peer_col_pos(&peer.alias, &peer.cols, c).is_some_and(|p| {
matches!(
peer.cols[p].ty,
spg_storage::DataType::Int
| spg_storage::DataType::BigInt
| spg_storage::DataType::SmallInt
)
}),
Expr::Literal(spg_sql::ast::Literal::Integer(_)) => true,
Expr::Binary { lhs, op, rhs } => {
matches!(op, BinOp::Add | BinOp::Sub | BinOp::Mul)
&& int_only_key_expr(lhs, peer)
&& int_only_key_expr(rhs, peer)
}
_ => false,
}
}
fn int_only_left_expr(e: &Expr, left_slice: &[ColumnSchema]) -> bool {
use spg_sql::ast::BinOp;
match e {
Expr::Column(c) => Engine::composite_col_pos(left_slice, c).is_some_and(|p| {
matches!(
left_slice[p].ty,
spg_storage::DataType::Int
| spg_storage::DataType::BigInt
| spg_storage::DataType::SmallInt
)
}),
Expr::Literal(spg_sql::ast::Literal::Integer(_)) => true,
Expr::Binary { lhs, op, rhs } => {
matches!(op, BinOp::Add | BinOp::Sub | BinOp::Mul)
&& int_only_left_expr(lhs, left_slice)
&& int_only_left_expr(rhs, left_slice)
}
_ => false,
}
}
fn eval_int_only_probe(
e: &Expr,
left_slice: &[ColumnSchema],
sources: &[JoinSrc<'_>],
offsets: &[usize],
tuple: &[usize],
) -> Result<Option<i64>, EngineError> {
use spg_sql::ast::BinOp;
match e {
Expr::Column(c) => {
let pos = Engine::composite_col_pos(left_slice, c).expect("classifier-checked");
Ok(match tuple_value(sources, offsets, tuple, pos) {
Some(Value::BigInt(n)) => Some(*n),
Some(Value::Int(n)) => Some(i64::from(*n)),
Some(Value::SmallInt(n)) => Some(i64::from(*n)),
_ => None,
})
}
Expr::Literal(spg_sql::ast::Literal::Integer(n)) => Ok(Some(*n)),
Expr::Binary { lhs, op, rhs } => {
let (Some(a), Some(b)) = (
eval_int_only_probe(lhs, left_slice, sources, offsets, tuple)?,
eval_int_only_probe(rhs, left_slice, sources, offsets, tuple)?,
) else {
return Ok(None);
};
let out = match op {
BinOp::Add => a.checked_add(b),
BinOp::Sub => a.checked_sub(b),
BinOp::Mul => a.checked_mul(b),
_ => unreachable!("classifier admits Add/Sub/Mul only"),
};
out.map(Some).ok_or_else(|| {
EngineError::Eval(crate::eval::EvalError::TypeMismatch {
detail: "bigint out of range".into(),
})
})
}
_ => unreachable!("classifier admits columns/integers/arithmetic only"),
}
}
fn expr_mentions_a_column(e: &Expr) -> bool {
match e {
Expr::Column(_) => true,
Expr::Unary { expr, .. } | Expr::Cast { expr, .. } => expr_mentions_a_column(expr),
Expr::Binary { lhs, rhs, .. } => expr_mentions_a_column(lhs) || expr_mentions_a_column(rhs),
_ => false,
}
}
fn match_equi_pair(
sub: &Expr,
peer: &JoinedPeer<'_>,
combined_schema: &[ColumnSchema],
consumed_cols: usize,
) -> Option<(usize, usize)> {
let Expr::Binary {
lhs,
op: spg_sql::ast::BinOp::Eq,
rhs,
} = sub
else {
return None;
};
let (Expr::Column(a), Expr::Column(b)) = (lhs.as_ref(), rhs.as_ref()) else {
return None;
};
let left_slice = &combined_schema[..consumed_cols];
if let (Some(l), Some(r)) = (
Engine::composite_col_pos(left_slice, a),
Engine::peer_col_pos(&peer.alias, &peer.cols, b),
) {
return Some((l, r));
}
if let (Some(l), Some(r)) = (
Engine::composite_col_pos(left_slice, b),
Engine::peer_col_pos(&peer.alias, &peer.cols, a),
) {
return Some((l, r));
}
None
}
fn where_equi_candidates<'w>(from: &FromClause, where_: Option<&'w Expr>) -> Vec<&'w Expr> {
let Some(w) = where_ else { return Vec::new() };
if !from
.joins
.iter()
.all(|j| matches!(j.kind, JoinKind::Inner | JoinKind::Cross))
{
return Vec::new();
}
reorder::split_and_conjunctions(w)
.into_iter()
.filter(|sub| {
matches!(
sub,
Expr::Binary { lhs, op: spg_sql::ast::BinOp::Eq, rhs }
if matches!((lhs.as_ref(), rhs.as_ref()), (Expr::Column(_), Expr::Column(_)))
)
})
.collect()
}
impl Engine {
#[allow(clippy::type_complexity)]
fn build_join_peers<'a>(
&self,
from: &'a FromClause,
peer_preds: &[Vec<&Expr>],
needed: Option<&alloc::collections::BTreeSet<(String, String)>>,
budget: &mut ByteBudget,
) -> Result<Vec<JoinedPeer<'a>>, EngineError> {
let mut joined: Vec<JoinedPeer<'a>> = Vec::new();
for j in &from.joins {
let a = j
.table
.alias
.as_deref()
.unwrap_or(j.table.name.as_str())
.to_string();
if let Some(inner_box) = &j.table.lateral_subquery {
if is_constant_values_derived(inner_box)
|| (derived_is_plain_table_select(inner_box, self.active_catalog())
&& !crate::subquery::select_is_correlated(inner_box))
{
let pidx = from
.joins
.iter()
.position(|jj| core::ptr::eq(jj, j))
.unwrap_or(0);
let (mut rows, mut cols) =
self.materialise_table_ref_filtered(&j.table, &peer_preds[pidx])?;
for (i, new_name) in j.table.unnest_column_aliases.iter().enumerate() {
if let Some(col) = cols.get_mut(i) {
col.name = new_name.clone();
}
}
if let Some(needed) = needed {
Self::null_out_unreferenced(&mut rows, &cols, &a, needed);
}
budget.charge(approx_rows_bytes(&rows))?;
joined.push(JoinedPeer {
eager_rows: Some(rows),
cols,
alias: a,
kind: j.kind,
on: j.on.as_ref(),
lateral: None,
join_table: None,
where_preds: Vec::new(),
});
continue;
}
let mut schema = self.lateral_probe_schema(inner_box)?;
for (i, new_name) in j.table.unnest_column_aliases.iter().enumerate() {
if let Some(col) = schema.get_mut(i) {
col.name = new_name.clone();
}
}
joined.push(JoinedPeer {
eager_rows: None,
cols: schema,
alias: a,
kind: j.kind,
on: j.on.as_ref(),
lateral: Some(inner_box.as_ref()),
join_table: None,
where_preds: Vec::new(),
});
} else {
let pidx = from
.joins
.iter()
.position(|jj| core::ptr::eq(jj, j))
.unwrap_or(0);
let plain = j.table.unnest_expr.is_none() && j.table.as_of_segment.is_none();
if plain && let Some(t) = self.active_catalog().get(&j.table.name) {
const SMALL_PEER_EAGER_ROWS: usize = 256;
let has_pushdown = !peer_preds[pidx].is_empty();
let peer_total = t.rows().len();
if has_pushdown && peer_total <= SMALL_PEER_EAGER_ROWS {
let (mut rows, cols) =
self.materialise_table_ref_filtered(&j.table, &peer_preds[pidx])?;
if let Some(needed) = needed {
Self::null_out_unreferenced(&mut rows, &cols, &a, needed);
}
budget.charge(approx_rows_bytes(&rows))?;
joined.push(JoinedPeer {
eager_rows: Some(rows),
cols,
alias: a,
kind: j.kind,
on: j.on.as_ref(),
lateral: None,
join_table: Some(j.table.name.clone()),
where_preds: Vec::new(),
});
continue;
}
joined.push(JoinedPeer {
eager_rows: None,
cols: t.schema().columns.clone(),
alias: a,
kind: j.kind,
on: j.on.as_ref(),
lateral: None,
join_table: Some(j.table.name.clone()),
where_preds: peer_preds[pidx].iter().map(|e| (*e).clone()).collect(),
});
continue;
}
let (mut rows, cols) =
self.materialise_table_ref_filtered(&j.table, &peer_preds[pidx])?;
if let Some(needed) = needed {
Self::null_out_unreferenced(&mut rows, &cols, &a, needed);
}
budget.charge(approx_rows_bytes(&rows))?;
joined.push(JoinedPeer {
eager_rows: Some(rows),
cols,
alias: a,
kind: j.kind,
on: j.on.as_ref(),
lateral: None,
join_table: Some(j.table.name.clone()),
where_preds: Vec::new(),
});
}
}
Ok(joined)
}
pub(crate) fn build_joined_filtered_rows(
&self,
from: &FromClause,
where_: Option<&Expr>,
cancel: CancelToken<'_>,
needed: Option<&alloc::collections::BTreeSet<(String, String)>>,
budget: &mut ByteBudget,
) -> Result<DeferredJoin<'_>, EngineError> {
let (swapped_from, primary_preds, peer_preds) = analyze_join_pushdown(from, where_);
let mut pushed_set: alloc::collections::BTreeSet<usize> = primary_preds
.iter()
.chain(peer_preds.iter().flat_map(|v| v.iter()))
.map(|e| core::ptr::from_ref::<Expr>(*e) as usize)
.collect();
let from = swapped_from.as_ref().unwrap_or(from);
let primary_alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str())
.to_string();
let primary_table: Option<&Table> = if !from.joins.is_empty()
&& from.primary.unnest_expr.is_none()
&& from.primary.lateral_subquery.is_none()
&& from.primary.as_of_segment.is_none()
{
self.active_catalog().get(&from.primary.name).filter(|t|
!t.has_cold_rows_fast())
} else {
None
};
let (primary_rows, primary_cols, primary_indices) = match primary_table {
Some(t) => {
let idxs = self.filter_table_indices(t, &primary_alias, &primary_preds)?;
let scan_snapshot = self.current_snapshot();
let idxs: Vec<usize> = idxs
.into_iter()
.filter(|&i| t.is_row_visible(i, &scan_snapshot))
.collect();
(Vec::new(), t.schema().columns.clone(), Some(idxs))
}
None => {
let (mut rows, cols) =
self.materialise_table_ref_filtered(&from.primary, &primary_preds)?;
if let Some(needed) = needed {
Self::null_out_unreferenced(&mut rows, &cols, &primary_alias, needed);
}
budget.charge(approx_rows_bytes(&rows))?;
(rows, cols, None)
}
};
let mut joined = self.build_join_peers(from, &peer_preds, needed, budget)?;
let combined_schema = build_combined_schema(&primary_alias, &primary_cols, &joined);
let join_sess = self.dml_session();
let ctx = EvalContext::new(&combined_schema, None)
.with_catalog(self.active_catalog())
.with_session(&join_sess);
if joined.is_empty() {
let mut filtered: Vec<Row<'static>> = Vec::new();
let mut memo = memoize::MemoizeCache::default();
for row in primary_rows {
if let Some(where_expr) = where_ {
let cond = self.eval_expr_with_correlated(
where_expr,
&row,
&ctx,
cancel,
Some(&mut memo),
)?;
if !crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)? {
continue;
}
}
filtered.push(row);
}
let width = combined_schema.len();
let n = filtered.len();
let offsets = alloc::vec![0, width];
let pos_to_src = build_pos_to_src(&offsets);
return Ok(DeferredJoin {
sources: alloc::vec![JoinSrc::Owned(filtered)],
offsets,
pos_to_src,
widths: alloc::vec![width],
masks: alloc::vec![None],
survivors: (0..n).collect(),
stride: 1,
combined_schema,
});
}
let primary_width = primary_cols.len();
#[allow(clippy::type_complexity)]
let (primary_source, primary_mask, working): (
JoinSrc<'_>,
Option<Vec<bool>>,
Vec<usize>,
) = match primary_indices {
Some(idxs) => {
let t = primary_table.expect("stored primary");
(
JoinSrc::Stored(t.rows()),
keep_mask(needed, &primary_cols, &primary_alias),
idxs,
)
}
None => {
let n = primary_rows.len();
(JoinSrc::Owned(primary_rows), None, (0..n).collect())
}
};
let where_equi = where_equi_candidates(from, where_);
let mut pipe = JoinPipeline::new(primary_source, primary_mask, primary_width, working);
for peer in &mut joined {
if pipe.rows() > MAX_JOIN_INTERMEDIATE_ROWS {
return Err(EngineError::Unsupported(alloc::format!(
"join intermediate result exceeds {MAX_JOIN_INTERMEDIATE_ROWS} rows ({} so far) - add join predicates",
pipe.rows()
)));
}
let right_arity = peer.cols.len();
let peer_mask = keep_mask(needed, &peer.cols, &peer.alias);
let (mut eq_pairs, eq_exprs, eq_probe_exprs, residual) =
extract_join_keys(peer, &combined_schema, pipe.consumed_cols);
if peer.lateral.is_none() && matches!(peer.kind, JoinKind::Inner | JoinKind::Cross) {
for cand in &where_equi {
if let Some(pair) =
match_equi_pair(cand, peer, &combined_schema, pipe.consumed_cols)
&& !eq_pairs.contains(&pair)
{
eq_pairs.push(pair);
pushed_set.insert(core::ptr::from_ref::<Expr>(*cand) as usize);
}
}
}
let extra_preds = core::mem::take(&mut peer.where_preds);
let residual: Vec<&Expr> = residual.into_iter().chain(extra_preds.iter()).collect();
if !matches!(peer.kind, JoinKind::Semi)
&& self.join_stage_inl(
&mut pipe,
peer,
&eq_pairs,
&residual,
&peer_mask,
right_arity,
&ctx,
cancel,
)?
{
continue;
}
if (!eq_pairs.is_empty() || !eq_exprs.is_empty() || !eq_probe_exprs.is_empty())
&& peer.lateral.is_none()
{
self.join_stage_hash(
&mut pipe,
peer,
&eq_pairs,
&eq_exprs,
&eq_probe_exprs,
&residual,
&peer_mask,
right_arity,
&combined_schema,
&ctx,
cancel,
)?;
continue;
}
self.join_stage_nested(
&mut pipe,
peer,
right_arity,
&combined_schema,
&ctx,
cancel,
needed,
budget,
)?;
}
let residual_where_owned: Option<Expr> = where_.and_then(|w| {
let kept: Vec<Expr> = reorder::split_and_conjunctions(w)
.into_iter()
.filter(|c| !pushed_set.contains(&(core::ptr::from_ref::<Expr>(c) as usize)))
.cloned()
.collect();
kept.into_iter().reduce(|a, b| Expr::Binary {
lhs: alloc::boxed::Box::new(a),
op: spg_sql::ast::BinOp::And,
rhs: alloc::boxed::Box::new(b),
})
});
let survivors =
self.filter_join_survivors(&pipe, residual_where_owned.as_ref(), &ctx, cancel, budget)?;
Ok(DeferredJoin {
sources: pipe.sources,
offsets: pipe.offsets,
pos_to_src: pipe.pos_to_src,
widths: pipe.widths,
masks: pipe.masks,
survivors,
stride: pipe.stride,
combined_schema,
})
}
#[allow(clippy::too_many_arguments)]
fn join_stage_inl<'a, 'p>(
&'a self,
pipe: &mut JoinPipeline<'a>,
peer: &JoinedPeer<'p>,
eq_pairs: &[(usize, usize)],
residual: &[&Expr],
peer_mask: &Option<Vec<bool>>,
right_arity: usize,
ctx: &EvalContext,
cancel: CancelToken<'_>,
) -> Result<bool, EngineError> {
const INL_MAX_LEFT: usize = 1024;
if matches!(peer.kind, JoinKind::Right | JoinKind::FullOuter) {
return Ok(false);
}
let Some(tname) = &peer.join_table else {
return Ok(false);
};
if !(peer.eager_rows.is_none() && !eq_pairs.is_empty() && pipe.rows() <= INL_MAX_LEFT) {
return Ok(false);
}
let Some(table) = self.active_catalog().get(tname) else {
return Ok(false);
};
let Some(idx) = peer
.cols
.iter()
.position(|c| c.name == peer.cols[eq_pairs[0].1].name)
.and_then(|pos| table.index_on(pos))
else {
return Ok(false);
};
let has_cold = table.has_cold_rows_fast();
let pk_col_pos = table
.schema()
.uniqueness_constraints
.iter()
.find(|u| u.is_primary_key && u.columns.len() == 1)
.map(|u| u.columns[0]);
let join_col_is_pk = pk_col_pos == Some(idx.column_position);
if has_cold && !join_col_is_pk {
return Ok(false);
}
let (cold_rows, cold_pk_map): (Vec<Row<'static>>, hashbrown::HashMap<i64, usize>) =
if has_cold {
crate::constraints::iter_cold_rows_with_locator_map(self.active_catalog(), table)
} else {
(Vec::new(), hashbrown::HashMap::new())
};
let stored = table.rows();
let hot_len = stored.len();
let scan_snapshot = self.current_snapshot();
let (lpos0, _) = eq_pairs[0];
let mut next: Vec<usize> = Vec::new();
for tuple in pipe.working.chunks(pipe.stride) {
cancel.check()?;
let mut left_matched = false;
if let Some(kv) = tuple_value(&pipe.sources, &pipe.offsets, tuple, lpos0)
&& !matches!(kv, Value::Null)
&& let Some(key) = spg_storage::IndexKey::from_value(kv)
{
for loc in idx.lookup_eq(&key) {
let ri = match *loc {
spg_storage::RowLocator::Hot(i) => i,
spg_storage::RowLocator::Cold { .. } => {
let spg_storage::IndexKey::Int(pk) = &key else {
continue;
};
match cold_pk_map.get(pk) {
Some(&off) => hot_len + off,
None => continue,
}
}
};
let right_opt: Option<&Row<'static>> = if ri < hot_len {
if !table.is_row_visible(ri, &scan_snapshot) {
continue;
}
stored.get(ri)
} else {
cold_rows.get(ri - hot_len)
};
let right = match right_opt {
Some(r) => r,
None => continue,
};
let mut ok = true;
for (lp, rp) in eq_pairs.iter().skip(1) {
let lv = tuple_value(&pipe.sources, &pipe.offsets, tuple, *lp);
let rv = right.values.get(*rp);
let eq = match (lv, rv) {
(Some(a), Some(b)) => {
!matches!(a, Value::Null)
&& !matches!(b, Value::Null)
&& value_cmp(a, b) == core::cmp::Ordering::Equal
}
_ => false,
};
if !eq {
ok = false;
break;
}
}
if !ok {
continue;
}
let keep = if residual.is_empty() {
true
} else {
let mut combined_vals = materialise_tuple_vals(
&pipe.sources,
&pipe.widths,
&pipe.masks,
tuple,
pipe.consumed_cols + right_arity,
);
extend_masked(&mut combined_vals, right, peer_mask.as_deref());
let combined = Row::new(combined_vals);
let mut k = true;
for r in residual {
let cond =
self.eval_expr_with_correlated(r, &combined, ctx, cancel, None)?;
if !crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)?
{
k = false;
break;
}
}
k
};
if keep {
next.extend_from_slice(tuple);
next.push(ri);
left_matched = true;
}
}
}
if !left_matched && matches!(peer.kind, JoinKind::Left) {
next.extend_from_slice(tuple);
next.push(usize::MAX);
}
}
let src = if cold_rows.is_empty() {
JoinSrc::Stored(stored)
} else {
JoinSrc::Mixed {
hot: stored,
cold: cold_rows,
cold_locator_map: cold_pk_map,
}
};
pipe.advance(next, src, peer_mask.clone(), right_arity);
Ok(true)
}
#[allow(clippy::too_many_arguments)]
fn join_stage_hash<'a, 'p>(
&'a self,
pipe: &mut JoinPipeline<'a>,
peer: &mut JoinedPeer<'p>,
eq_pairs: &[(usize, usize)],
eq_exprs: &[(usize, &Expr, &Expr)],
eq_probe_exprs: &[(usize, &Expr, &Expr)],
residual: &[&Expr],
peer_mask: &Option<Vec<bool>>,
right_arity: usize,
combined_schema: &[ColumnSchema],
ctx: &EvalContext,
cancel: CancelToken<'_>,
) -> Result<(), EngineError> {
let (rights_src, build_gate): (JoinSrc<'a>, Option<(&'a Table, usize)>) =
match peer.eager_rows.take() {
Some(rows) => (JoinSrc::Owned(rows), None),
None => match peer
.join_table
.as_deref()
.and_then(|n| self.active_catalog().get(n))
{
Some(t) if t.has_cold_rows_fast() => {
let (cold, map) = crate::constraints::iter_cold_rows_with_locator_map(
self.active_catalog(),
t,
);
let hot = t.rows();
let hot_len = hot.len();
(
JoinSrc::Mixed {
hot,
cold,
cold_locator_map: map,
},
Some((t, hot_len)),
)
}
Some(t) => (JoinSrc::Stored(t.rows()), Some((t, t.rows().len()))),
None => (JoinSrc::Owned(Vec::new()), None),
},
};
let scan_snapshot = self.current_snapshot();
let n_rights = rights_src.len();
let int_keyed = eq_exprs.is_empty()
&& eq_probe_exprs.is_empty()
&& eq_pairs.len() == 1
&& matches!(
combined_schema[eq_pairs[0].0].ty,
spg_storage::DataType::BigInt
| spg_storage::DataType::Int
| spg_storage::DataType::SmallInt
)
&& {
let peer_col_ty = peer.cols.get(eq_pairs[0].1).map(|c| c.ty);
matches!(
peer_col_ty,
Some(
spg_storage::DataType::BigInt
| spg_storage::DataType::Int
| spg_storage::DataType::SmallInt
)
)
};
let int_expr_keyed = eq_pairs.is_empty()
&& eq_probe_exprs.is_empty()
&& eq_exprs.len() == 1
&& matches!(
combined_schema[eq_exprs[0].0].ty,
spg_storage::DataType::BigInt
| spg_storage::DataType::Int
| spg_storage::DataType::SmallInt
)
&& int_only_key_expr(eq_exprs[0].1, peer);
let int_probe_expr_keyed =
eq_pairs.is_empty() && eq_exprs.is_empty() && eq_probe_exprs.len() == 1;
let int2_keyed = eq_exprs.is_empty()
&& eq_pairs.len() + eq_probe_exprs.len() == 2
&& !eq_probe_exprs.is_empty()
&& eq_pairs.iter().all(|(l, r)| {
matches!(
combined_schema[*l].ty,
spg_storage::DataType::BigInt
| spg_storage::DataType::Int
| spg_storage::DataType::SmallInt
) && matches!(
peer.cols.get(*r).map(|c| c.ty),
Some(
spg_storage::DataType::BigInt
| spg_storage::DataType::Int
| spg_storage::DataType::SmallInt
)
)
});
let residual: Vec<&Expr> = if int_expr_keyed {
residual
.iter()
.copied()
.filter(|r| !core::ptr::eq(*r, eq_exprs[0].2))
.collect()
} else if int_probe_expr_keyed {
residual
.iter()
.copied()
.filter(|r| !core::ptr::eq(*r, eq_probe_exprs[0].2))
.collect()
} else if !eq_probe_exprs.is_empty() {
residual
.iter()
.copied()
.filter(|r| !eq_probe_exprs.iter().any(|(_, _, c)| core::ptr::eq(*r, *c)))
.collect()
} else {
residual.to_vec()
};
let mut build_preds: Vec<eval::CompiledExpr> = Vec::new();
let residual: Vec<&Expr> = {
let peer_ctx_probe = EvalContext::new(&peer.cols, Some(peer.alias.as_str()));
residual
.iter()
.copied()
.filter(|r| {
let peer_only = {
let all = core::cell::Cell::new(true);
crate::expr_analysis::visit_expr_columns_and_subqueries(
r,
&mut |c| {
if Engine::peer_col_pos(&peer.alias, &peer.cols, c).is_none() {
all.set(false);
}
},
&mut |_| {
all.set(false);
},
);
all.get()
};
if peer_only && eval::fully_compilable(r) && expr_mentions_a_column(r) {
build_preds.push(eval::compile_expr(r, &peer_ctx_probe));
false
} else {
true
}
})
.collect()
};
let residual = residual.as_slice();
let any_int_lane = int_keyed || int_expr_keyed || int_probe_expr_keyed;
let mut int2_table: hashbrown::HashMap<i128, Bucket> =
hashbrown::HashMap::with_capacity(if int2_keyed { n_rights } else { 0 });
let mut table: hashbrown::HashMap<String, Bucket> =
hashbrown::HashMap::with_capacity(if any_int_lane { 0 } else { n_rights });
let mut int_table: hashbrown::HashMap<i64, Bucket> =
hashbrown::HashMap::with_capacity(if any_int_lane { n_rights } else { 0 });
let peer_ctx = EvalContext {
columns: &peer.cols,
table_alias: Some(peer.alias.as_str()),
..ctx.clone()
};
let mut keybuf: Vec<&Value> = Vec::with_capacity(eq_pairs.len());
let mut pred_stack: Vec<Value<'static>> = Vec::new();
let mut keystr = String::new();
let mut built_parallel = false;
if (any_int_lane || int2_keyed)
&& n_rights >= crate::PARALLEL_MIN_ROWS
&& !matches!(rights_src, JoinSrc::Mixed { .. })
&& let Some(r) = self.parallel_runner.0.as_deref()
{
struct ShardTables {
t64: hashbrown::HashMap<i64, Bucket>,
t128: hashbrown::HashMap<i128, Bucket>,
}
type ShardOut = Result<ShardTables, EngineError>;
let n_shards = (n_rights / crate::PARALLEL_MIN_ROWS).clamp(2, 8);
let chunk = n_rights.div_ceil(n_shards);
let rights_ref = &rights_src;
let preds_ref = &build_preds;
let peer_cols = &peer.cols;
let peer_alias_s = peer.alias.as_str();
let mysql = ctx.mysql_dialect;
let style = ctx.render_style;
let cat = ctx.catalog;
let eq_pairs_ref = eq_pairs;
let eq_exprs_ref = eq_exprs;
let eq_probe_ref = eq_probe_exprs;
let results = r.run_shards(n_shards, &|si| {
let lo = si * chunk;
let hi = ((si + 1) * chunk).min(n_rights);
let mut sctx = EvalContext::new(peer_cols, Some(peer_alias_s));
sctx.mysql_dialect = mysql;
sctx.render_style = style;
let sctx = match cat {
Some(c) => sctx.with_catalog(c),
None => sctx,
};
let mut stack: Vec<Value<'static>> = Vec::new();
let mut out = ShardTables {
t64: hashbrown::HashMap::new(),
t128: hashbrown::HashMap::new(),
};
let run = || -> ShardOut {
let mut out = out;
'srows: for ri in lo..hi {
if let Some((gt, hot_len)) = build_gate
&& ri < hot_len
&& !gt.is_row_visible(ri, &scan_snapshot)
{
continue;
}
let Some(right) = rights_ref.get(ri) else {
continue;
};
for c in preds_ref.iter() {
let v = eval::eval_compiled(c, right, &sctx, &mut stack)
.map_err(EngineError::Eval)?;
if !crate::eval::predicate_is_true(&v, "JOIN/ON", mysql)? {
continue 'srows;
}
}
if int2_keyed {
let mut parts = [0i64; 2];
let mut pi = 0;
for (_, rpos) in eq_pairs_ref {
match right.values.get(*rpos) {
Some(Value::BigInt(n)) => parts[pi] = *n,
Some(Value::Int(n)) => parts[pi] = i64::from(*n),
Some(Value::SmallInt(n)) => parts[pi] = i64::from(*n),
_ => continue 'srows,
}
pi += 1;
}
for (p, _, _) in eq_probe_ref {
match right.values.get(*p) {
Some(Value::BigInt(n)) => parts[pi] = *n,
Some(Value::Int(n)) => parts[pi] = i64::from(*n),
Some(Value::SmallInt(n)) => parts[pi] = i64::from(*n),
_ => continue 'srows,
}
pi += 1;
}
let key = ((parts[0] as i128) << 64) | (parts[1] as u64 as i128);
match out.t128.entry(key) {
hashbrown::hash_map::Entry::Occupied(mut o) => o.get_mut().push(ri),
hashbrown::hash_map::Entry::Vacant(v) => {
v.insert(Bucket::One(ri));
}
}
continue;
}
let key: i64 = if !eq_pairs_ref.is_empty() || !eq_probe_ref.is_empty() {
let rpos = if eq_probe_ref.is_empty() {
eq_pairs_ref[0].1
} else {
eq_probe_ref[0].0
};
match right.values.get(rpos) {
Some(Value::BigInt(n)) => *n,
Some(Value::Int(n)) => i64::from(*n),
Some(Value::SmallInt(n)) => i64::from(*n),
_ => continue 'srows,
}
} else {
match eval::eval_expr(eq_exprs_ref[0].1, right, &sctx)
.map_err(EngineError::Eval)?
{
Value::BigInt(n) => n,
Value::Int(n) => i64::from(n),
Value::SmallInt(n) => i64::from(n),
_ => continue 'srows,
}
};
match out.t64.entry(key) {
hashbrown::hash_map::Entry::Occupied(mut o) => o.get_mut().push(ri),
hashbrown::hash_map::Entry::Vacant(v) => {
v.insert(Bucket::One(ri));
}
}
}
Ok(out)
};
alloc::boxed::Box::new(run())
});
let mut ok = true;
let mut shard_tables: Vec<ShardTables> = Vec::with_capacity(n_shards);
let mut first_err: Option<EngineError> = None;
for boxed in results {
match boxed.downcast::<ShardOut>() {
Ok(sh) => match *sh {
Ok(t) => shard_tables.push(t),
Err(e) => {
ok = false;
if first_err.is_none() {
first_err = Some(e);
}
}
},
Err(_) => ok = false,
}
}
if let Some(e) = first_err {
return Err(e);
}
if ok {
for t in shard_tables {
for (k, b) in t.t64 {
match int_table.entry(k) {
hashbrown::hash_map::Entry::Occupied(mut o) => {
for ri in b.as_slice() {
o.get_mut().push(*ri);
}
}
hashbrown::hash_map::Entry::Vacant(v) => {
v.insert(b);
}
}
}
for (k, b) in t.t128 {
match int2_table.entry(k) {
hashbrown::hash_map::Entry::Occupied(mut o) => {
for ri in b.as_slice() {
o.get_mut().push(*ri);
}
}
hashbrown::hash_map::Entry::Vacant(v) => {
v.insert(b);
}
}
}
}
built_parallel = true;
}
}
'build: for ri in 0..n_rights {
if built_parallel {
break;
}
if let Some((gt, hot_len)) = build_gate
&& ri < hot_len
&& !gt.is_row_visible(ri, &scan_snapshot)
{
continue;
}
let Some(right) = rights_src.get(ri) else {
continue;
};
if !build_preds.is_empty() {
let mut keep = true;
for c in &build_preds {
let v = eval::eval_compiled(c, right, &peer_ctx, &mut pred_stack)
.map_err(EngineError::Eval)?;
if !crate::eval::predicate_is_true(&v, "JOIN/ON", ctx.mysql_dialect)? {
keep = false;
break;
}
}
if !keep {
continue 'build;
}
}
if int2_keyed {
let mut parts = [0i64; 2];
let mut pi = 0;
let mut null_key = false;
for (_, rpos) in eq_pairs {
match right.values.get(*rpos) {
Some(Value::BigInt(n)) => parts[pi] = *n,
Some(Value::Int(n)) => parts[pi] = i64::from(*n),
Some(Value::SmallInt(n)) => parts[pi] = i64::from(*n),
_ => {
null_key = true;
break;
}
}
pi += 1;
}
if !null_key {
for (p, _, _) in eq_probe_exprs {
match right.values.get(*p) {
Some(Value::BigInt(n)) => parts[pi] = *n,
Some(Value::Int(n)) => parts[pi] = i64::from(*n),
Some(Value::SmallInt(n)) => parts[pi] = i64::from(*n),
_ => {
null_key = true;
break;
}
}
pi += 1;
}
}
if null_key {
continue 'build;
}
let key = ((parts[0] as i128) << 64) | (parts[1] as u64 as i128);
match int2_table.entry(key) {
hashbrown::hash_map::Entry::Occupied(mut o) => o.get_mut().push(ri),
hashbrown::hash_map::Entry::Vacant(v) => {
v.insert(Bucket::One(ri));
}
}
continue;
}
if any_int_lane {
let key = if int_keyed || int_probe_expr_keyed {
let rpos = if int_keyed {
eq_pairs[0].1
} else {
eq_probe_exprs[0].0
};
match right.values.get(rpos) {
Some(Value::BigInt(n)) => *n,
Some(Value::Int(n)) => i64::from(*n),
Some(Value::SmallInt(n)) => i64::from(*n),
_ => continue 'build,
}
} else {
match eval::eval_expr(eq_exprs[0].1, right, &peer_ctx)
.map_err(EngineError::Eval)?
{
Value::BigInt(n) => n,
Value::Int(n) => i64::from(n),
Value::SmallInt(n) => i64::from(n),
_ => continue 'build,
}
};
match int_table.entry(key) {
hashbrown::hash_map::Entry::Occupied(mut o) => o.get_mut().push(ri),
hashbrown::hash_map::Entry::Vacant(v) => {
v.insert(Bucket::One(ri));
}
}
continue;
}
keybuf.clear();
for (_, rpos) in eq_pairs {
match right.values.get(*rpos) {
Some(v) if !matches!(v, Value::Null) => keybuf.push(v),
_ => continue 'build,
}
}
aggregate::encode_key_refs_into(&keybuf, &mut keystr);
for (_, e, _) in eq_exprs {
let v = eval::eval_expr(e, right, &peer_ctx).map_err(EngineError::Eval)?;
if matches!(v, Value::Null) {
continue 'build;
}
aggregate::push_canonical_key(&mut keystr, &v);
}
for (p, _, _) in eq_probe_exprs {
match right.values.get(*p) {
Some(v) if !matches!(v, Value::Null) => {
aggregate::push_canonical_key(&mut keystr, v);
}
_ => continue 'build,
}
}
match table.get_mut(keystr.as_str()) {
Some(b) => b.push(ri),
None => {
table.insert(keystr.clone(), Bucket::One(ri));
}
}
}
let mut next: Vec<usize> = Vec::new();
let track_right = matches!(peer.kind, JoinKind::Right | JoinKind::FullOuter);
let mut peer_matched: Vec<bool> = if track_right {
alloc::vec![false; n_rights]
} else {
Vec::new()
};
let mut probebuf: Vec<&Value> = Vec::with_capacity(eq_pairs.len());
for tuple in pipe.working.chunks(pipe.stride) {
cancel.check()?;
let mut left_matched = false;
let mut left_has_null = false;
let int2_probe_key: Option<i128> = if int2_keyed {
let mut parts = [0i64; 2];
let mut pi = 0;
let mut nul = false;
for (lpos, _) in eq_pairs {
match tuple_value(&pipe.sources, &pipe.offsets, tuple, *lpos) {
Some(Value::BigInt(n)) => parts[pi] = *n,
Some(Value::Int(n)) => parts[pi] = i64::from(*n),
Some(Value::SmallInt(n)) => parts[pi] = i64::from(*n),
_ => {
nul = true;
break;
}
}
pi += 1;
}
if !nul {
for (_, e, _) in eq_probe_exprs {
match eval_int_only_probe(
e,
&combined_schema[..pipe.consumed_cols],
&pipe.sources,
&pipe.offsets,
tuple,
)? {
Some(k) => parts[pi] = k,
None => {
nul = true;
break;
}
}
pi += 1;
}
}
if nul {
left_has_null = true;
None
} else {
Some(((parts[0] as i128) << 64) | (parts[1] as u64 as i128))
}
} else {
None
};
let int_probe_key: Option<i64> = if int2_keyed {
None
} else if int_probe_expr_keyed {
match eval_int_only_probe(
eq_probe_exprs[0].1,
&combined_schema[..pipe.consumed_cols],
&pipe.sources,
&pipe.offsets,
tuple,
)? {
Some(k) => Some(k),
None => {
left_has_null = true;
None
}
}
} else if any_int_lane {
let lpos = if int_keyed {
eq_pairs[0].0
} else {
eq_exprs[0].0
};
match tuple_value(&pipe.sources, &pipe.offsets, tuple, lpos) {
Some(Value::BigInt(n)) => Some(*n),
Some(Value::Int(n)) => Some(i64::from(*n)),
Some(Value::SmallInt(n)) => Some(i64::from(*n)),
_ => {
left_has_null = true;
None
}
}
} else {
probebuf.clear();
for (lpos, _) in eq_pairs {
match tuple_value(&pipe.sources, &pipe.offsets, tuple, *lpos) {
Some(v) if !matches!(v, Value::Null) => probebuf.push(v),
_ => {
left_has_null = true;
break;
}
}
}
if !left_has_null {
aggregate::encode_key_refs_into(&probebuf, &mut keystr);
for (lpos, _, _) in eq_exprs {
match tuple_value(&pipe.sources, &pipe.offsets, tuple, *lpos) {
Some(v) if !matches!(v, Value::Null) => {
aggregate::push_canonical_key(&mut keystr, v)
}
_ => {
left_has_null = true;
break;
}
}
}
}
if !left_has_null {
for (_, e, _) in eq_probe_exprs {
match eval_int_only_probe(
e,
&combined_schema[..pipe.consumed_cols],
&pipe.sources,
&pipe.offsets,
tuple,
)? {
Some(k) => {
aggregate::push_canonical_key(&mut keystr, &Value::BigInt(k))
}
None => {
left_has_null = true;
break;
}
}
}
}
None
};
let cands_opt: Option<&Bucket> = if left_has_null {
None
} else if int2_keyed {
int2_table.get(&int2_probe_key.unwrap())
} else if any_int_lane {
int_table.get(&int_probe_key.unwrap())
} else {
table.get(keystr.as_str())
};
if let Some(cands) = cands_opt {
for &ri in cands.as_slice() {
let keep = if residual.is_empty() {
true
} else {
let right = rights_src.get(ri).expect("hash candidate row");
let mut combined_vals = materialise_tuple_vals(
&pipe.sources,
&pipe.widths,
&pipe.masks,
tuple,
pipe.consumed_cols + right_arity,
);
extend_masked(&mut combined_vals, right, peer_mask.as_deref());
let combined = Row::new(combined_vals);
let mut ok = true;
for r in residual {
let cond =
self.eval_expr_with_correlated(r, &combined, ctx, cancel, None)?;
if !crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)?
{
ok = false;
break;
}
}
ok
};
if keep {
next.extend_from_slice(tuple);
next.push(ri);
left_matched = true;
if track_right {
peer_matched[ri] = true;
}
if matches!(peer.kind, JoinKind::Semi) {
break;
}
}
}
}
if !left_matched && matches!(peer.kind, JoinKind::Left | JoinKind::FullOuter) {
next.extend_from_slice(tuple);
next.push(usize::MAX);
}
}
if track_right {
for ri in 0..n_rights {
if peer_matched[ri] {
continue;
}
if let Some((gt, hot_len)) = build_gate
&& ri < hot_len
&& !gt.is_row_visible(ri, &scan_snapshot)
{
continue;
}
if rights_src.get(ri).is_none() {
continue;
}
for _ in 0..pipe.stride {
next.push(usize::MAX);
}
next.push(ri);
}
}
pipe.advance(next, rights_src, peer_mask.clone(), right_arity);
debug_assert!(pipe.consumed_cols <= combined_schema.len());
Ok(())
}
#[allow(clippy::too_many_arguments)]
fn join_stage_nested<'a, 'p>(
&'a self,
pipe: &mut JoinPipeline<'a>,
peer: &mut JoinedPeer<'p>,
right_arity: usize,
combined_schema: &[ColumnSchema],
ctx: &EvalContext,
cancel: CancelToken<'_>,
needed: Option<&alloc::collections::BTreeSet<(String, String)>>,
budget: &mut ByteBudget,
) -> Result<(), EngineError> {
let lazy_rows: Option<Vec<Row<'static>>> =
if peer.eager_rows.is_none() && peer.lateral.is_none() {
let tname = peer.join_table.as_deref().unwrap_or("");
let snap = self.current_snapshot();
let mut rows: Vec<Row<'static>> = self
.active_catalog()
.get(tname)
.map(|t| t.scan_visible(&snap).map(|(_, r)| r.clone()).collect())
.unwrap_or_default();
if let Some(t) = self.active_catalog().get(tname)
&& t.has_cold_rows_fast()
{
rows.extend(crate::constraints::iter_cold_rows_of_parent(
self.active_catalog(),
t,
));
}
if let Some(needed) = needed {
Self::null_out_unreferenced(&mut rows, &peer.cols, &peer.alias, needed);
}
budget.charge(approx_rows_bytes(&rows))?;
Some(rows)
} else {
None
};
let mut arena: Vec<Row<'static>> = Vec::new();
let rights_eager: Option<&[Row<'static>]> =
peer.eager_rows.as_deref().or(lazy_rows.as_deref());
let mut next: Vec<usize> = Vec::new();
let right_or_full = matches!(peer.kind, JoinKind::Right | JoinKind::FullOuter);
let lateral_fixed: Option<Vec<Row<'static>>> =
if right_or_full && let Some(inner) = peer.lateral {
match self.exec_select_cancel(inner, cancel)? {
QueryResult::Rows { rows, .. } => Some(rows),
_ => {
return Err(EngineError::Unsupported(
"derived-table join operand must be a SELECT".into(),
));
}
}
} else {
None
};
let track_right = right_or_full && (peer.lateral.is_none() || lateral_fixed.is_some());
let fixed_peer_len = match &lateral_fixed {
Some(f) => f.len(),
None => rights_eager.map(<[_]>::len).unwrap_or(0),
};
let mut peer_matched: Vec<bool> = if track_right {
alloc::vec![false; fixed_peer_len]
} else {
Vec::new()
};
for tuple in pipe.working.chunks(pipe.stride) {
cancel.check()?;
let mut left_matched = false;
let left_vals = materialise_tuple_vals(
&pipe.sources,
&pipe.widths,
&pipe.masks,
tuple,
pipe.consumed_cols,
);
let per_left_rrows: Cow<'_, [Row]> = match (&lateral_fixed, peer.lateral) {
(Some(fixed), _) => Cow::Borrowed(fixed.as_slice()),
(None, Some(inner)) => {
let outer_schema = &combined_schema[..pipe.consumed_cols];
let left_row = Row::new(left_vals.clone());
let rows =
self.materialise_lateral_for_outer(inner, outer_schema, &left_row)?;
Cow::Owned(rows)
}
(None, None) => Cow::Borrowed(rights_eager.expect("non-lateral peer eager")),
};
for (ri, right) in per_left_rrows.as_ref().iter().enumerate() {
let mut combined_vals = left_vals.clone();
combined_vals.extend(right.values.iter().cloned());
let combined = Row::new(combined_vals);
let keep = if let Some(on_expr) = peer.on {
let cond =
self.eval_expr_with_correlated(on_expr, &combined, ctx, cancel, None)?;
crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)?
} else {
true
};
if keep {
next.extend_from_slice(tuple);
if peer.lateral.is_some() && lateral_fixed.is_none() {
let mut cv = combined.values;
let rv = cv.split_off(left_vals.len());
arena.push(Row::new(rv));
next.push(arena.len() - 1);
} else {
next.push(ri);
if track_right {
peer_matched[ri] = true;
}
}
left_matched = true;
if matches!(peer.kind, JoinKind::Semi) {
break;
}
}
}
if !left_matched && matches!(peer.kind, JoinKind::Left | JoinKind::FullOuter) {
next.extend_from_slice(tuple);
next.push(usize::MAX);
}
}
if track_right {
for (ri, matched) in peer_matched.iter().enumerate() {
if *matched {
continue;
}
for _ in 0..pipe.stride {
next.push(usize::MAX);
}
next.push(ri);
}
}
if next.len() / (pipe.stride + 1) > MAX_JOIN_INTERMEDIATE_ROWS {
return Err(EngineError::Unsupported(alloc::format!(
"join intermediate result exceeds {MAX_JOIN_INTERMEDIATE_ROWS} rows ({} so far) - add join predicates",
next.len() / (pipe.stride + 1)
)));
}
let source = if let Some(fixed) = lateral_fixed {
JoinSrc::Owned(fixed)
} else if peer.lateral.is_some() {
JoinSrc::Owned(arena)
} else if let Some(lz) = lazy_rows {
JoinSrc::Owned(lz)
} else {
JoinSrc::Owned(peer.eager_rows.take().expect("non-lateral peer eager"))
};
pipe.advance(next, source, None, right_arity);
debug_assert!(pipe.consumed_cols <= combined_schema.len());
Ok(())
}
fn filter_join_survivors(
&self,
pipe: &JoinPipeline<'_>,
where_: Option<&Expr>,
ctx: &EvalContext,
cancel: CancelToken<'_>,
budget: &mut ByteBudget,
) -> Result<Vec<usize>, EngineError> {
if where_.is_none() {
let n_rows = if pipe.stride == 0 {
0
} else {
pipe.working.len() / pipe.stride
};
if n_rows > 0 {
let sample_tuple = &pipe.working[..pipe.stride];
let per_tuple =
approx_tuple_bytes(&pipe.sources, &pipe.offsets, &pipe.masks, sample_tuple);
budget.charge(per_tuple.saturating_mul(n_rows))?;
}
cancel.check()?;
return Ok(pipe.working.clone());
}
let mut memo = memoize::MemoizeCache::default();
let compiled_where: Option<eval::CompiledExpr> = where_
.filter(|w| eval::fully_compilable(w))
.map(|w| eval::compile_expr(w, ctx));
let mut survivors: Vec<usize> = Vec::new();
for tuple in pipe.working.chunks(pipe.stride) {
let rr = RowRef::Tuple {
sources: &pipe.sources,
offsets: &pipe.offsets,
pos_to_src: &pipe.pos_to_src,
tuple,
};
let mut eval_stack: Vec<Value<'_>> = Vec::new();
let pass = if let Some(cw) = &compiled_where {
matches!(
eval::eval_compiled_ref(cw, rr, ctx, &mut eval_stack)
.map_err(EngineError::Eval)?,
Value::Bool(true)
)
} else if let Some(where_expr) = where_ {
let row = rr.as_row();
matches!(
self.eval_expr_with_correlated(where_expr, &row, ctx, cancel, Some(&mut memo))?,
Value::Bool(true)
)
} else {
true
};
if !pass {
continue;
}
budget.charge(approx_tuple_bytes(
&pipe.sources,
&pipe.offsets,
&pipe.masks,
tuple,
))?;
survivors.extend_from_slice(tuple);
}
Ok(survivors)
}
fn lateral_probe_schema(
&self,
inner: &SelectStatement,
) -> Result<Vec<ColumnSchema>, EngineError> {
match self.execute_readonly_select_for_lateral_probe(inner) {
Ok(QueryResult::Rows { columns, .. }) => Ok(columns),
_ => {
if let [SelectItem::Wildcard] = inner.items.as_slice()
&& let Some(from) = &inner.from
&& from.joins.is_empty()
&& (from.primary.unnest_expr.is_some()
|| from.primary.generate_series_args.is_some())
{
let t = &from.primary;
let elem_dtype = if t.generate_series_args.is_some() {
DataType::BigInt
} else {
DataType::Text
};
let first = t
.unnest_column_aliases
.first()
.cloned()
.or_else(|| t.alias.clone())
.unwrap_or_else(|| t.name.clone());
let mut out = alloc::vec![ColumnSchema::new(first, elem_dtype, true)];
if t.with_ordinality {
let ord = t
.unnest_column_aliases
.get(1)
.cloned()
.unwrap_or_else(|| "ordinality".to_string());
out.push(ColumnSchema::new(ord, DataType::BigInt, false));
}
return Ok(out);
}
if let Some(from) = &inner.from
&& from.joins.is_empty()
&& matches!(inner.items.as_slice(), [SelectItem::Wildcard])
&& let Some(jt) = from.primary.json_table.as_deref()
{
return Ok(crate::select::json_table_schema_pub(&jt.columns));
}
if let Some(from) = &inner.from
&& from.joins.is_empty()
&& matches!(inner.items.as_slice(), [SelectItem::Wildcard])
&& let Some((fn_name, _)) = from.primary.table_fn_call.as_deref()
{
let cat = self.active_catalog();
let overloads = cat.functions_named(fn_name);
if let Some(def) = overloads.first() {
let declared = def.returns.trim();
let upper = declared.to_ascii_uppercase();
if let Some(rest) = upper.strip_prefix("TABLE(") {
let _ = rest;
let raw = &declared["TABLE(".len()..declared.len() - 1];
let cols: Vec<ColumnSchema> = raw
.split(',')
.map(|decl| {
let cname = decl.split_whitespace().next().unwrap_or("col");
ColumnSchema::new(cname.to_string(), DataType::Text, true)
})
.collect();
return Ok(cols);
}
let cname = from
.primary
.alias
.clone()
.unwrap_or_else(|| fn_name.clone());
return Ok(alloc::vec![ColumnSchema::new(cname, DataType::Text, true)]);
}
}
let mut out: Vec<ColumnSchema> = Vec::new();
for (i, item) in inner.items.iter().enumerate() {
let name = match item {
SelectItem::Expr { alias: Some(a), .. } => a.clone(),
SelectItem::Expr { expr, .. } => synth_lateral_col_name(expr, i),
SelectItem::Wildcard | SelectItem::QualifiedWildcard(_) => {
alloc::format!("col{i}")
}
};
out.push(ColumnSchema::new(name, DataType::Text, true));
}
Ok(out)
}
}
}
fn execute_readonly_select_for_lateral_probe(
&self,
inner: &SelectStatement,
) -> Result<QueryResult, EngineError> {
self.exec_bare_select_cancel(inner, CancelToken::none())
}
fn materialise_lateral_for_outer(
&self,
inner: &SelectStatement,
outer_schema: &[ColumnSchema],
outer_row: &Row<'static>,
) -> Result<Vec<Row<'static>>, EngineError> {
let mut substituted = inner.clone();
substitute_outer_columns_multi(&mut substituted, outer_row, outer_schema);
let result = self.exec_bare_select_cancel(&substituted, CancelToken::none())?;
match result {
QueryResult::Rows { rows, .. } => Ok(rows),
_ => Err(EngineError::Unsupported(
"LATERAL subquery must be a SELECT (cannot be a write statement)".into(),
)),
}
}
pub(crate) fn try_count_star_left_anti_join_fast(
&self,
stmt: &SelectStatement,
from: &FromClause,
) -> Result<Option<QueryResult>, EngineError> {
use spg_sql::ast::{JoinKind, SelectItem};
ANTI_JOIN_FAST_PATH_TRIED.fetch_add(1, core::sync::atomic::Ordering::Relaxed);
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()
{
return Ok(None);
}
if from.joins.len() != 1 {
return Ok(None);
}
let join = &from.joins[0];
if !matches!(join.kind, JoinKind::Left) {
return Ok(None);
}
let plain = |t: &spg_sql::ast::TableRef| {
t.unnest_expr.is_none()
&& t.lateral_subquery.is_none()
&& t.as_of_segment.is_none()
&& t.generate_series_args.is_none()
};
if !plain(&from.primary) || !plain(&join.table) {
return Ok(None);
}
let outer_alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str());
let inner_alias = join
.table
.alias
.as_deref()
.unwrap_or(join.table.name.as_str());
if 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(on) = join.on.as_ref() else {
return Ok(None);
};
enum InnerKey {
Col(String),
Expr(Expr),
}
let (outer_col, inner_key, null_cols): (String, InnerKey, Vec<String>) =
if let Some((oc, ic)) = analyse_join_eq(on, outer_alias, inner_alias)? {
let nulls = alloc::vec![ic.clone()];
(oc, InnerKey::Col(ic), nulls)
} else if let Some((oc, e)) = analyse_join_eq_expr(on, outer_alias, inner_alias) {
let mut cols: Vec<String> = Vec::new();
collect_inner_int_cols(&e, &mut cols);
(oc, InnerKey::Expr(e), cols)
} else {
return Ok(None);
};
let Some(where_expr) = stmt.where_.as_ref() else {
return Ok(None);
};
if !null_cols
.iter()
.any(|c| is_inner_is_null(where_expr, inner_alias, c))
{
return Ok(None);
}
let catalog = self.active_catalog();
let Some(inner_table) = catalog.get(join.table.name.as_str()) else {
return Ok(None);
};
let inner_schema = inner_table.schema();
let int_col_pos = |name: &str| -> Option<usize> {
inner_schema
.columns
.iter()
.position(|c| c.name.eq_ignore_ascii_case(name))
.filter(|&p| {
matches!(
inner_schema.columns[p].ty,
spg_storage::DataType::BigInt
| spg_storage::DataType::Int
| spg_storage::DataType::SmallInt
)
})
};
let inner_pos: Option<usize> = match &inner_key {
InnerKey::Col(c) => {
let Some(p) = int_col_pos(c) else {
return Ok(None);
};
Some(p)
}
InnerKey::Expr(_) => {
if !null_cols.iter().all(|c| int_col_pos(c).is_some()) {
return Ok(None);
}
None
}
};
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(|c| c.name.eq_ignore_ascii_case(&outer_col))
else {
return Ok(None);
};
let outer_ty = outer_schema.columns[outer_pos].ty;
if !matches!(
outer_ty,
spg_storage::DataType::BigInt
| spg_storage::DataType::Int
| spg_storage::DataType::SmallInt
) {
return Ok(None);
}
let read_int = |v: &Value| -> Option<i64> {
match v {
Value::BigInt(n) => Some(*n),
Value::Int(n) => Some(i64::from(*n)),
Value::SmallInt(n) => Some(i64::from(*n)),
_ => None,
}
};
let scan_snapshot = self.current_snapshot();
let mut antiset: hashbrown::HashSet<i64> =
hashbrown::HashSet::with_capacity(inner_table.row_count());
let inner_ctx = self.ev_ctx(&inner_schema.columns, Some(inner_alias));
for (i, row) in inner_table.rows().iter().enumerate() {
if !inner_table.is_row_visible(i, &scan_snapshot) {
continue;
}
match (&inner_key, inner_pos) {
(InnerKey::Col(_), Some(p)) => {
if let Some(v) = row.values.get(p)
&& let Some(k) = read_int(v)
{
antiset.insert(k);
}
}
(InnerKey::Expr(e), _) => {
let v = eval::eval_expr(e, row, &inner_ctx).map_err(EngineError::Eval)?;
if let Some(k) = read_int(&v) {
antiset.insert(k);
}
}
_ => unreachable!("Col always carries a position"),
}
}
let mut count: i64 = 0;
for (i, row) in outer_table.rows().iter().enumerate() {
if !outer_table.is_row_visible(i, &scan_snapshot) {
continue;
}
match row.values.get(outer_pos) {
Some(v) => match read_int(v) {
Some(k) => {
if !antiset.contains(&k) {
count += 1;
}
}
None => count += 1,
},
None => 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 _ = outer_alias;
let _ = outer_col;
ANTI_JOIN_FAST_PATH_FIRED.fetch_add(1, core::sync::atomic::Ordering::Relaxed);
Ok(Some(QueryResult::Rows { columns, rows }))
}
pub(crate) fn try_streamed_inner_join_walk_topn(
&self,
stmt: &SelectStatement,
from: &FromClause,
cancel: CancelToken<'_>,
) -> Result<Option<QueryResult>, EngineError> {
let Some(limit) = stmt.limit_literal() else {
return Ok(None);
};
if stmt.offset.is_some() && stmt.offset_literal().is_none() {
return Ok(None);
}
if stmt.distinct
|| stmt.limit_with_ties
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| aggregate::uses_aggregate(stmt)
{
return Ok(None);
}
if from.joins.len() != 1 {
return Ok(None);
}
let j = &from.joins[0];
if !matches!(j.kind, JoinKind::Inner) {
return Ok(None);
}
let plain = |t: &TableRef| {
t.unnest_expr.is_none() && t.lateral_subquery.is_none() && t.as_of_segment.is_none()
};
if !plain(&from.primary) || !plain(&j.table) {
return Ok(None);
}
let Some(on_expr) = j.on.as_ref() else {
return Ok(None);
};
let Some(primary_table) = self.active_catalog().get(&from.primary.name) else {
return Ok(None);
};
if self.active_catalog().get(&j.table.name).is_none() {
return Ok(None);
}
let primary_alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str())
.to_string();
if stmt.order_by.len() != 1 {
return Ok(None);
}
let order = &stmt.order_by[0];
let Expr::Column(order_col) = &order.expr else {
return Ok(None);
};
if let Some(q) = &order_col.qualifier
&& !q.eq_ignore_ascii_case(&primary_alias)
{
return Ok(None);
}
let primary_cols = primary_table.schema().columns.clone();
let Some(order_col_pos) = primary_cols
.iter()
.position(|c| c.name.eq_ignore_ascii_case(&order_col.name))
else {
return Ok(None);
};
let Some(order_index) = primary_table.index_on(order_col_pos) else {
return Ok(None);
};
if !matches!(order_index.kind, spg_storage::IndexKind::BTree(_)) {
return Ok(None);
}
let peer_alias = j
.table
.alias
.as_deref()
.unwrap_or(j.table.name.as_str())
.to_string();
let mut needed = alloc::collections::BTreeSet::new();
let prunable = collect_qualified_refs(stmt, &mut needed).is_some();
let mut budget = ByteBudget::new(self.max_query_bytes);
let (mut peer_rows, peer_cols) = self.materialise_table_ref_filtered(&j.table, &[])?;
if prunable {
Self::null_out_unreferenced(&mut peer_rows, &peer_cols, &peer_alias, &needed);
}
budget.charge(approx_rows_bytes(&peer_rows))?;
let mut combined_schema: Vec<ColumnSchema> = Vec::new();
for col in &primary_cols {
combined_schema.push(ColumnSchema::new(
alloc::format!("{primary_alias}.{}", col.name),
col.ty,
col.nullable,
));
}
for col in &peer_cols {
combined_schema.push(ColumnSchema::new(
alloc::format!("{peer_alias}.{}", col.name),
col.ty,
col.nullable,
));
}
let join_sess = self.dml_session();
let ctx = EvalContext::new(&combined_schema, None)
.with_catalog(self.active_catalog())
.with_session(&join_sess);
let left_arity = primary_cols.len();
let mut eq_pairs: Vec<(usize, usize)> = Vec::new();
let mut residual: Vec<&Expr> = Vec::new();
for sub in reorder::split_and_conjunctions(on_expr) {
let mut matched = None;
if let Expr::Binary {
lhs,
op: spg_sql::ast::BinOp::Eq,
rhs,
} = sub
&& let (Expr::Column(a), Expr::Column(b)) = (lhs.as_ref(), rhs.as_ref())
{
let left_slice = &combined_schema[..left_arity];
if let (Some(l), Some(r)) = (
Self::composite_col_pos(left_slice, a),
Self::peer_col_pos(&peer_alias, &peer_cols, b),
) {
matched = Some((l, r));
} else if let (Some(l), Some(r)) = (
Self::composite_col_pos(left_slice, b),
Self::peer_col_pos(&peer_alias, &peer_cols, a),
) {
matched = Some((l, r));
}
}
match matched {
Some(pair) => eq_pairs.push(pair),
None => residual.push(sub),
}
}
if eq_pairs.is_empty() {
return Ok(None);
}
let mut htable: hashbrown::HashMap<String, Vec<usize>> =
hashbrown::HashMap::with_capacity(peer_rows.len());
let mut keybuf: Vec<Value<'static>> = Vec::with_capacity(eq_pairs.len());
'build: for (ri, right) in peer_rows.iter().enumerate() {
keybuf.clear();
for (_, rpos) in &eq_pairs {
let v = right.values.get(*rpos).cloned().unwrap_or(Value::Null);
if matches!(v, Value::Null) {
continue 'build;
}
keybuf.push(v);
}
htable
.entry(aggregate::encode_key(&keybuf))
.or_default()
.push(ri);
}
let keep_mask: Vec<bool> = primary_cols
.iter()
.map(|c| !prunable || needed.contains(&(primary_alias.clone(), c.name.clone())))
.collect();
let keep = (limit as usize).saturating_add(stmt.offset_literal().map_or(0, |o| o as usize));
let mut where_memo = memoize::MemoizeCache::default();
let mut plain_sink: Vec<Row<'static>> = Vec::with_capacity(keep.min(1024));
let outer_schema: &[ColumnSchema] = &combined_schema[..left_arity];
let outer_ctx = EvalContext::new(outer_schema, None).with_catalog(self.active_catalog());
let where_conjuncts: Vec<&Expr> = stmt
.where_
.as_ref()
.map(|w| reorder::split_and_conjunctions(w))
.unwrap_or_default();
let (where_outer_only, where_mixed): (Vec<&Expr>, Vec<&Expr>) =
where_conjuncts.iter().copied().partition(|e| {
crate::joinfold::expr_references_alias(e, &primary_alias)
&& !crate::joinfold::expr_references_any_other_alias(e, &primary_alias)
});
let (residual_outer_only, residual_mixed): (Vec<&Expr>, Vec<&Expr>) =
residual.iter().copied().partition(|e| {
crate::joinfold::expr_references_alias(e, &primary_alias)
&& !crate::joinfold::expr_references_any_other_alias(e, &primary_alias)
});
let mut outer_memo = memoize::MemoizeCache::default();
let walker: alloc::boxed::Box<
dyn Iterator<Item = (&spg_storage::IndexKey, &Vec<spg_storage::RowLocator>)>,
> = if order.desc {
alloc::boxed::Box::new(order_index.iter_desc())
} else {
alloc::boxed::Box::new(order_index.iter_asc())
};
let primary_table_name = primary_table.schema().name.clone();
let scan_snapshot = self.current_snapshot();
'walk: for (key, locators) in walker {
cancel.check()?;
for loc in locators {
let left_cow: Cow<'_, Row> = match *loc {
spg_storage::RowLocator::Hot(i) => {
if !primary_table.is_row_visible(i, &scan_snapshot) {
continue;
}
match primary_table.rows().get(i) {
Some(r) => Cow::Borrowed(r),
None => continue,
}
}
spg_storage::RowLocator::Cold { segment_id, .. } => {
match self.active_catalog().resolve_cold_locator(
&primary_table_name,
segment_id,
key,
) {
Some(r) => Cow::Owned(r),
None => continue,
}
}
};
let left: &Row<'static> = left_cow.as_ref();
keybuf.clear();
let mut left_has_null = false;
for (lpos, _) in &eq_pairs {
let v = left.values.get(*lpos).cloned().unwrap_or(Value::Null);
if matches!(v, Value::Null) {
left_has_null = true;
break;
}
keybuf.push(v);
}
if left_has_null {
continue;
}
let Some(cands) = htable.get(&aggregate::encode_key(&keybuf)) else {
continue;
};
let mut outer_ok = true;
for r in &residual_outer_only {
let cond = self.eval_expr_with_correlated(r, left, &outer_ctx, cancel, None)?;
if !crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)? {
outer_ok = false;
break;
}
}
if !outer_ok {
continue;
}
for w in &where_outer_only {
let cond = self.eval_expr_with_correlated(
w,
left,
&outer_ctx,
cancel,
Some(&mut outer_memo),
)?;
if !crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)? {
outer_ok = false;
break;
}
}
if !outer_ok {
continue;
}
for &ri in cands {
let right = &peer_rows[ri];
let mut combined_vals: Vec<Value<'static>> =
Vec::with_capacity(left_arity + peer_cols.len());
for (i, v) in left.values.iter().enumerate() {
combined_vals.push(if keep_mask.get(i).copied().unwrap_or(true) {
v.clone()
} else {
Value::Null
});
}
combined_vals.extend(right.values.iter().cloned());
let combined = Row::new(combined_vals);
let mut ok = true;
for r in &residual_mixed {
let cond =
self.eval_expr_with_correlated(r, &combined, &ctx, cancel, None)?;
if !crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)? {
ok = false;
break;
}
}
if !ok {
continue;
}
for w in &where_mixed {
let cond = self.eval_expr_with_correlated(
w,
&combined,
&ctx,
cancel,
Some(&mut where_memo),
)?;
if !crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)? {
ok = false;
break;
}
}
if !ok {
continue;
}
budget.charge(approx_row_bytes(&combined))?;
plain_sink.push(combined);
if plain_sink.len() >= keep {
break 'walk;
}
}
}
}
let mut output = plain_sink;
apply_offset_and_limit(&mut output, stmt.offset_literal(), stmt.limit_literal());
let projection =
build_projection(&stmt.items, &combined_schema, "", self.backslash_escapes)?;
let mut proj_memo = memoize::MemoizeCache::default();
let mut rows: Vec<Row<'static>> = Vec::with_capacity(output.len());
for row in &output {
let mut values = Vec::with_capacity(projection.len());
for p in &projection {
values.push(self.eval_expr_with_correlated(
&p.expr,
row,
&ctx,
cancel,
Some(&mut proj_memo),
)?);
}
rows.push(Row::new(values));
}
let columns: Vec<ColumnSchema> = projection
.into_iter()
.map(|p| ColumnSchema::new(p.output_name, p.ty, p.nullable))
.collect();
Ok(Some(QueryResult::Rows { columns, rows }))
}
pub(crate) fn try_streamed_inner_join_topn(
&self,
stmt: &SelectStatement,
from: &FromClause,
cancel: CancelToken<'_>,
) -> Result<Option<QueryResult>, EngineError> {
let Some(limit) = stmt.limit_literal() else {
return Ok(None);
};
if stmt.offset.is_some() && stmt.offset_literal().is_none() {
return Ok(None);
}
if stmt.distinct
|| stmt.group_by.is_some()
|| stmt.having.is_some()
|| aggregate::uses_aggregate(stmt)
{
return Ok(None);
}
if from.joins.len() != 1 {
return Ok(None);
}
let j = &from.joins[0];
if !matches!(j.kind, JoinKind::Inner) {
return Ok(None);
}
let plain = |t: &TableRef| {
t.unnest_expr.is_none() && t.lateral_subquery.is_none() && t.as_of_segment.is_none()
};
if !plain(&from.primary) || !plain(&j.table) {
return Ok(None);
}
let Some(on_expr) = j.on.as_ref() else {
return Ok(None);
};
let Some(primary_table) = self.active_catalog().get(&from.primary.name) else {
return Ok(None);
};
if self.active_catalog().get(&j.table.name).is_none() {
return Ok(None);
}
let primary_alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str())
.to_string();
let peer_alias = j
.table
.alias
.as_deref()
.unwrap_or(j.table.name.as_str())
.to_string();
let mut needed = alloc::collections::BTreeSet::new();
let prunable = collect_qualified_refs(stmt, &mut needed).is_some();
let mut budget = ByteBudget::new(self.max_query_bytes);
let (mut peer_rows, peer_cols) = self.materialise_table_ref_filtered(&j.table, &[])?;
if prunable {
Self::null_out_unreferenced(&mut peer_rows, &peer_cols, &peer_alias, &needed);
}
budget.charge(approx_rows_bytes(&peer_rows))?;
let primary_cols = primary_table.schema().columns.clone();
let mut combined_schema: Vec<ColumnSchema> = Vec::new();
for col in &primary_cols {
combined_schema.push(ColumnSchema::new(
alloc::format!("{primary_alias}.{}", col.name),
col.ty,
col.nullable,
));
}
for col in &peer_cols {
combined_schema.push(ColumnSchema::new(
alloc::format!("{peer_alias}.{}", col.name),
col.ty,
col.nullable,
));
}
let join_sess = self.dml_session();
let ctx = EvalContext::new(&combined_schema, None)
.with_catalog(self.active_catalog())
.with_session(&join_sess);
let left_arity = primary_cols.len();
let mut eq_pairs: Vec<(usize, usize)> = Vec::new();
let mut residual: Vec<&Expr> = Vec::new();
for sub in reorder::split_and_conjunctions(on_expr) {
let mut matched = None;
if let Expr::Binary {
lhs,
op: spg_sql::ast::BinOp::Eq,
rhs,
} = sub
&& let (Expr::Column(a), Expr::Column(b)) = (lhs.as_ref(), rhs.as_ref())
{
let left_slice = &combined_schema[..left_arity];
if let (Some(l), Some(r)) = (
Self::composite_col_pos(left_slice, a),
Self::peer_col_pos(&peer_alias, &peer_cols, b),
) {
matched = Some((l, r));
} else if let (Some(l), Some(r)) = (
Self::composite_col_pos(left_slice, b),
Self::peer_col_pos(&peer_alias, &peer_cols, a),
) {
matched = Some((l, r));
}
}
match matched {
Some(pair) => eq_pairs.push(pair),
None => residual.push(sub),
}
}
if eq_pairs.is_empty() {
return Ok(None); }
let mut htable: hashbrown::HashMap<String, Vec<usize>> =
hashbrown::HashMap::with_capacity(peer_rows.len());
let mut keybuf: Vec<Value<'static>> = Vec::with_capacity(eq_pairs.len());
'build: for (ri, right) in peer_rows.iter().enumerate() {
keybuf.clear();
for (_, rpos) in &eq_pairs {
let v = right.values.get(*rpos).cloned().unwrap_or(Value::Null);
if matches!(v, Value::Null) {
continue 'build;
}
keybuf.push(v);
}
htable
.entry(aggregate::encode_key(&keybuf))
.or_default()
.push(ri);
}
let keep_mask: Vec<bool> = primary_cols
.iter()
.map(|c| !prunable || needed.contains(&(primary_alias.clone(), c.name.clone())))
.collect();
let keep = (limit as usize).saturating_add(stmt.offset_literal().map_or(0, |o| o as usize));
let descs: alloc::rc::Rc<[bool]> = stmt
.order_by
.iter()
.map(|o| o.desc)
.collect::<Vec<bool>>()
.into();
let mut where_memo = memoize::MemoizeCache::default();
let mut heap: alloc::collections::BinaryHeap<TopNEntry> =
alloc::collections::BinaryHeap::new();
let mut plain_sink: Vec<Row<'static>> = Vec::new();
let mut seq: u64 = 0;
let primary_cold = self.iter_cold_rows_of_table(primary_table);
let snap = self.current_snapshot();
'scan: for left in primary_table
.scan_visible(&snap)
.map(|(_, r)| r)
.chain(primary_cold.iter())
{
cancel.check()?;
if keep == 0 {
break 'scan;
}
keybuf.clear();
let mut left_has_null = false;
for (lpos, _) in &eq_pairs {
let v = left.values.get(*lpos).cloned().unwrap_or(Value::Null);
if matches!(v, Value::Null) {
left_has_null = true;
break;
}
keybuf.push(v);
}
if left_has_null {
continue;
}
let Some(cands) = htable.get(&aggregate::encode_key(&keybuf)) else {
continue;
};
for &ri in cands {
let right = &peer_rows[ri];
let mut combined_vals: Vec<Value<'static>> =
Vec::with_capacity(left_arity + peer_cols.len());
for (i, v) in left.values.iter().enumerate() {
combined_vals.push(if keep_mask.get(i).copied().unwrap_or(true) {
v.clone()
} else {
Value::Null
});
}
combined_vals.extend(right.values.iter().cloned());
let combined = Row::new(combined_vals);
let mut ok = true;
for r in &residual {
let cond = self.eval_expr_with_correlated(r, &combined, &ctx, cancel, None)?;
if !crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)? {
ok = false;
break;
}
}
if !ok {
continue;
}
if let Some(w) = stmt.where_.as_ref() {
let cond = self.eval_expr_with_correlated(
w,
&combined,
&ctx,
cancel,
Some(&mut where_memo),
)?;
if !crate::eval::predicate_is_true(&cond, "JOIN/ON", ctx.mysql_dialect)? {
continue;
}
}
if stmt.order_by.is_empty() {
budget.charge(approx_row_bytes(&combined))?;
plain_sink.push(combined);
if plain_sink.len() >= keep {
break 'scan;
}
} else {
let keys = build_order_keys(&stmt.order_by, &combined, &ctx)?;
let entry = TopNEntry {
keys,
descs: alloc::rc::Rc::clone(&descs),
seq,
row: combined,
};
seq += 1;
if heap.len() < keep {
budget.charge(approx_row_bytes(&entry.row))?;
heap.push(entry);
} else if let Some(top) = heap.peek()
&& entry < *top
{
if let Some(evicted) = heap.pop() {
budget.release(approx_row_bytes(&evicted.row));
}
budget.charge(approx_row_bytes(&entry.row))?;
heap.push(entry);
}
}
}
}
let mut output: Vec<Row<'static>> = if stmt.order_by.is_empty() {
plain_sink
} else {
heap.into_sorted_vec().into_iter().map(|e| e.row).collect()
};
apply_offset_and_limit(&mut output, stmt.offset_literal(), stmt.limit_literal());
let projection =
build_projection(&stmt.items, &combined_schema, "", self.backslash_escapes)?;
let mut proj_memo = memoize::MemoizeCache::default();
let mut rows: Vec<Row<'static>> = Vec::with_capacity(output.len());
for row in &output {
let mut values = Vec::with_capacity(projection.len());
for p in &projection {
values.push(self.eval_expr_with_correlated(
&p.expr,
row,
&ctx,
cancel,
Some(&mut proj_memo),
)?);
}
rows.push(Row::new(values));
}
let columns: Vec<ColumnSchema> = projection
.into_iter()
.map(|p| ColumnSchema::new(p.output_name, p.ty, p.nullable))
.collect();
Ok(Some(QueryResult::Rows { columns, rows }))
}
}
pub(crate) fn synth_lateral_col_name(expr: &Expr, idx: usize) -> String {
match expr {
Expr::Column(c) => c.name.clone(),
Expr::FunctionCall { name, .. } => name.clone(),
Expr::Cast { expr: inner, .. } => synth_lateral_col_name(inner, idx),
_ => alloc::format!("column{}", idx + 1),
}
}
fn expr_is_constant(e: &Expr) -> bool {
match e {
Expr::Literal(_) | Expr::Placeholder(_) => true,
Expr::Unary { expr, .. } | Expr::Cast { expr, .. } => expr_is_constant(expr),
Expr::Binary { lhs, rhs, .. } => expr_is_constant(lhs) && expr_is_constant(rhs),
_ => false,
}
}
fn derived_is_plain_table_select(s: &SelectStatement, cat: &crate::Catalog) -> bool {
let Some(from) = &s.from else {
return false;
};
let plain = |t: &spg_sql::ast::TableRef| {
t.unnest_expr.is_none()
&& t.generate_series_args.is_none()
&& t.lateral_subquery.is_none()
&& cat.get(&t.name).is_some()
};
plain(&from.primary) && from.joins.iter().all(|j| plain(&j.table))
}
fn is_constant_values_derived(s: &SelectStatement) -> bool {
use spg_sql::ast::SelectItem;
let peer_ok = |p: &SelectStatement| -> bool {
p.from.is_none()
&& p.where_.is_none()
&& p.having.is_none()
&& p.items.iter().all(|it| match it {
SelectItem::Expr { expr, .. } => expr_is_constant(expr),
SelectItem::Wildcard | SelectItem::QualifiedWildcard(_) => false,
})
};
peer_ok(s) && s.unions.iter().all(|(_, peer)| peer_ok(peer))
}
pub(crate) fn substitute_outer_columns_multi(
stmt: &mut SelectStatement,
outer_row: &Row<'static>,
outer_schema: &[ColumnSchema],
) {
substitute_outer_in_select(stmt, outer_row, outer_schema);
}
fn substitute_outer_in_select(
stmt: &mut SelectStatement,
outer_row: &Row<'static>,
outer_schema: &[ColumnSchema],
) {
let bare = stmt.from.is_none();
for item in &mut stmt.items {
if let SelectItem::Expr { expr, .. } = item {
substitute_outer_in_expr(expr, outer_row, outer_schema, bare);
}
}
if let Some(from) = &mut stmt.from {
substitute_outer_in_table_ref(&mut from.primary, outer_row, outer_schema);
for j in &mut from.joins {
substitute_outer_in_table_ref(&mut j.table, outer_row, outer_schema);
if let Some(on) = &mut j.on {
substitute_outer_in_expr(on, outer_row, outer_schema, bare);
}
}
}
if let Some(w) = &mut stmt.where_ {
substitute_outer_in_expr(w, outer_row, outer_schema, bare);
}
if let Some(gs) = &mut stmt.group_by {
for g in gs {
substitute_outer_in_expr(g, outer_row, outer_schema, bare);
}
}
if let Some(h) = &mut stmt.having {
substitute_outer_in_expr(h, outer_row, outer_schema, bare);
}
for o in &mut stmt.order_by {
substitute_outer_in_expr(&mut o.expr, outer_row, outer_schema, bare);
}
for (_, peer) in &mut stmt.unions {
substitute_outer_in_select(peer, outer_row, outer_schema);
}
}
fn substitute_outer_in_table_ref(
t: &mut spg_sql::ast::TableRef,
outer_row: &Row<'static>,
outer_schema: &[ColumnSchema],
) {
if let Some((_, arg)) = t.jsonb_each_text_arg.as_mut() {
substitute_outer_in_expr(arg, outer_row, outer_schema, true);
}
if let Some(arg) = t.unnest_expr.as_deref_mut() {
substitute_outer_in_expr(arg, outer_row, outer_schema, true);
}
if let Some(call) = t.table_fn_call.as_deref_mut() {
for a in call.1.iter_mut() {
substitute_outer_in_expr(a, outer_row, outer_schema, true);
}
}
if let Some(args) = t.generate_series_args.as_mut() {
for a in args.iter_mut() {
substitute_outer_in_expr(a, outer_row, outer_schema, true);
}
}
if let Some(inner) = t.lateral_subquery.as_deref_mut() {
substitute_outer_in_select(inner, outer_row, outer_schema);
}
if let Some(jt) = t.json_table.as_deref_mut() {
substitute_outer_in_expr(&mut jt.doc, outer_row, outer_schema, true);
for (_, e) in jt.passing.iter_mut() {
substitute_outer_in_expr(e, outer_row, outer_schema, true);
}
}
}
fn outer_col_index(
outer_schema: &[ColumnSchema],
qualifier: Option<&str>,
name: &str,
bare_ok: bool,
) -> Option<usize> {
match qualifier {
Some(q) => {
let composite = alloc::format!("{q}.{name}");
outer_schema
.iter()
.position(|sc| sc.name.eq_ignore_ascii_case(&composite))
}
None if bare_ok => {
let mut found = None;
for (i, sc) in outer_schema.iter().enumerate() {
let bare = sc.name.rsplit('.').next().unwrap_or(sc.name.as_str());
if bare.eq_ignore_ascii_case(name) {
if found.is_some() {
return None; }
found = Some(i);
}
}
found
}
None => None,
}
}
fn outer_value_to_expr(v: Value<'static>) -> Option<Expr> {
match v {
Value::TextArray(items) => Some(Expr::Array(
items
.into_iter()
.map(|it| {
Expr::Literal(match it {
Some(s) => spg_sql::ast::Literal::String(s),
None => spg_sql::ast::Literal::Null,
})
})
.collect(),
)),
Value::IntArray(items) => Some(Expr::Array(
items
.into_iter()
.map(|it| {
Expr::Literal(match it {
Some(n) => spg_sql::ast::Literal::Integer(i64::from(n)),
None => spg_sql::ast::Literal::Null,
})
})
.collect(),
)),
Value::BigIntArray(items) => Some(Expr::Array(
items
.into_iter()
.map(|it| {
Expr::Literal(match it {
Some(n) => spg_sql::ast::Literal::Integer(n),
None => spg_sql::ast::Literal::Null,
})
})
.collect(),
)),
other => value_to_literal_expr(other).ok(),
}
}
fn substitute_outer_in_expr(
e: &mut Expr,
outer_row: &Row<'static>,
outer_schema: &[ColumnSchema],
bare_ok: bool,
) {
if let Expr::Column(c) = e
&& let Some(idx) = outer_col_index(outer_schema, c.qualifier.as_deref(), &c.name, bare_ok)
{
let v = outer_row.values.get(idx).cloned().unwrap_or(Value::Null);
if let Some(lit) = outer_value_to_expr(v) {
*e = lit;
return;
}
}
let mut rec = |e: &mut Expr| substitute_outer_in_expr(e, outer_row, outer_schema, bare_ok);
match e {
Expr::Binary { lhs, rhs, .. } => {
rec(lhs);
rec(rhs);
}
Expr::Unary { expr: inner, .. } => rec(inner),
Expr::FunctionCall { args, .. } => {
for a in args {
rec(a);
}
}
Expr::Cast { expr: inner, .. } => rec(inner),
Expr::Array(items) => {
for it in items {
rec(it);
}
}
Expr::ArraySubscript { target, index } => {
rec(target);
rec(index);
}
Expr::ArraySlice { target, lo, hi } => {
rec(target);
if let Some(lo) = lo {
rec(lo);
}
if let Some(hi) = hi {
rec(hi);
}
}
Expr::InList { expr, list, .. } => {
rec(expr);
for it in list {
rec(it);
}
}
Expr::Case {
operand,
branches,
else_branch,
} => {
if let Some(op) = operand {
rec(op);
}
for (cond, val) in branches {
rec(cond);
rec(val);
}
if let Some(e) = else_branch {
rec(e);
}
}
_ => {}
}
}
fn analyze_join_pushdown<'w>(
from: &FromClause,
where_: Option<&'w Expr>,
) -> (Option<FromClause>, Vec<&'w Expr>, Vec<Vec<&'w Expr>>) {
let primary_alias = from
.primary
.alias
.as_deref()
.unwrap_or(from.primary.name.as_str());
let mut primary_preds: Vec<&Expr> = Vec::new();
let mut peer_preds: Vec<Vec<&Expr>> = alloc::vec![Vec::new(); from.joins.len()];
let primary_nullable = from
.joins
.iter()
.any(|j| matches!(j.kind, JoinKind::Right | JoinKind::FullOuter));
if let Some(w) = where_ {
for sub in reorder::split_and_conjunctions(w) {
if expr_has_subquery(sub) || aggregate::contains_aggregate(sub) {
continue;
}
let mut quals: Vec<&str> = Vec::new();
let mut all_qualified = true;
collect_column_qualifiers(sub, &mut quals, &mut all_qualified);
if !all_qualified || quals.is_empty() {
continue;
}
let q0 = quals[0];
if !quals.iter().all(|q| q.eq_ignore_ascii_case(q0)) {
continue;
}
if q0.eq_ignore_ascii_case(primary_alias) {
if !primary_nullable {
primary_preds.push(sub);
}
continue;
}
for (i, j) in from.joins.iter().enumerate() {
if matches!(j.kind, JoinKind::Inner | JoinKind::Cross)
&& j.table.lateral_subquery.is_none()
&& q0.eq_ignore_ascii_case(
j.table.alias.as_deref().unwrap_or(j.table.name.as_str()),
)
{
peer_preds[i].push(sub);
break;
}
}
}
}
if primary_preds.is_empty()
&& let Some(j0) = from.joins.first()
&& matches!(j0.kind, JoinKind::Inner)
&& j0.table.lateral_subquery.is_none()
&& !peer_preds[0].is_empty()
{
let peer_alias = j0.table.alias.as_deref().unwrap_or(j0.table.name.as_str());
let on_safe = j0.on.as_ref().is_some_and(|on| {
let mut quals: Vec<&str> = Vec::new();
let mut all_q = true;
collect_column_qualifiers(on, &mut quals, &mut all_q);
all_q
&& quals.iter().all(|q| {
q.eq_ignore_ascii_case(primary_alias) || q.eq_ignore_ascii_case(peer_alias)
})
});
if on_safe {
let mut from_owned = from.clone();
core::mem::swap(&mut from_owned.primary, &mut from_owned.joins[0].table);
let primary_preds = peer_preds[0].drain(..).collect();
return (Some(from_owned), primary_preds, peer_preds);
}
}
(None, primary_preds, peer_preds)
}
fn build_combined_schema(
primary_alias: &str,
primary_cols: &[ColumnSchema],
joined: &[JoinedPeer<'_>],
) -> Vec<ColumnSchema> {
let carry = |name: alloc::string::String, col: &ColumnSchema| {
let mut c = ColumnSchema::new(name, col.ty, col.nullable);
c.collation_name = col.collation_name.clone();
c.user_enum_type = col.user_enum_type.clone();
c
};
let mut combined_schema: Vec<ColumnSchema> = Vec::new();
for col in primary_cols {
combined_schema.push(carry(alloc::format!("{primary_alias}.{}", col.name), col));
}
for peer in joined {
for col in &peer.cols {
combined_schema.push(carry(alloc::format!("{}.{}", peer.alias, col.name), col));
}
}
combined_schema
}
fn analyse_join_eq(
on: &Expr,
outer_alias: &str,
inner_alias: &str,
) -> Result<Option<(String, String)>, EngineError> {
use spg_sql::ast::BinOp;
let Expr::Binary {
lhs,
op: BinOp::Eq,
rhs,
} = on
else {
return Ok(None);
};
let (Expr::Column(a), Expr::Column(b)) = (lhs.as_ref(), rhs.as_ref()) else {
return Ok(None);
};
fn col_alias(c: &spg_sql::ast::ColumnName) -> Option<&str> {
c.qualifier.as_deref()
}
let pair_o_then_i = (col_alias(a), col_alias(b));
if matches!(pair_o_then_i, (Some(aq), Some(bq))
if aq.eq_ignore_ascii_case(outer_alias) && bq.eq_ignore_ascii_case(inner_alias))
{
return Ok(Some((a.name.clone(), b.name.clone())));
}
if matches!(pair_o_then_i, (Some(aq), Some(bq))
if aq.eq_ignore_ascii_case(inner_alias) && bq.eq_ignore_ascii_case(outer_alias))
{
return Ok(Some((b.name.clone(), a.name.clone())));
}
Ok(None)
}
fn analyse_join_eq_expr(on: &Expr, outer_alias: &str, inner_alias: &str) -> Option<(String, Expr)> {
use spg_sql::ast::BinOp;
let Expr::Binary {
lhs,
op: BinOp::Eq,
rhs,
} = on
else {
return None;
};
fn inner_only_int(e: &Expr, inner_alias: &str) -> bool {
use spg_sql::ast::BinOp;
match e {
Expr::Column(c) => c
.qualifier
.as_deref()
.is_some_and(|q| q.eq_ignore_ascii_case(inner_alias)),
Expr::Literal(spg_sql::ast::Literal::Integer(_)) => true,
Expr::Binary { lhs, op, rhs } => {
matches!(op, BinOp::Add | BinOp::Sub | BinOp::Mul)
&& inner_only_int(lhs, inner_alias)
&& inner_only_int(rhs, inner_alias)
}
_ => false,
}
}
for (a, b) in [(lhs.as_ref(), rhs.as_ref()), (rhs.as_ref(), lhs.as_ref())] {
if let Expr::Column(c) = a
&& c.qualifier
.as_deref()
.is_some_and(|q| q.eq_ignore_ascii_case(outer_alias))
&& !matches!(b, Expr::Column(_))
&& inner_only_int(b, inner_alias)
&& expr_mentions_a_column(b)
{
return Some((c.name.clone(), b.clone()));
}
}
None
}
fn collect_inner_int_cols(e: &Expr, out: &mut Vec<String>) {
match e {
Expr::Column(c) => out.push(c.name.clone()),
Expr::Binary { lhs, rhs, .. } => {
collect_inner_int_cols(lhs, out);
collect_inner_int_cols(rhs, out);
}
_ => {}
}
}
fn is_inner_is_null(e: &Expr, inner_alias: &str, inner_col: &str) -> bool {
let Expr::IsNull { expr, negated } = e else {
return false;
};
if *negated {
return false;
}
let Expr::Column(c) = expr.as_ref() else {
return false;
};
c.qualifier
.as_deref()
.is_some_and(|q| q.eq_ignore_ascii_case(inner_alias))
&& c.name.eq_ignore_ascii_case(inner_col)
}
pub static ANTI_JOIN_FAST_PATH_TRIED: core::sync::atomic::AtomicU64 =
core::sync::atomic::AtomicU64::new(0);
pub static ANTI_JOIN_FAST_PATH_FIRED: core::sync::atomic::AtomicU64 =
core::sync::atomic::AtomicU64::new(0);
#[cfg(test)]
mod r655_rowref_size {
fn rowref_stays_small() {
assert_eq!(
core::mem::size_of::<super::RowRef<'_>>(),
64,
"RowRef changed size; a table scan allocates one per row, so \
this multiplies by the row count — re-measure scan RSS"
);
}
}