pub use super::dataset::{
convert_property_path, ConcreteStoreDataset, Dataset, DatasetPathAdapter, InMemoryDataset,
};
use crate::algebra::{Algebra, Binding, Solution, Term, Variable};
use anyhow::Result;
use std::collections::HashSet;
use super::queryexecutor_queries::substitute_algebra_binding;
use super::queryexecutor_type::QueryExecutor;
const BOUND_JOIN_ROW_THRESHOLD: usize = 100;
impl QueryExecutor {
pub(super) fn execute_index_optimized_join(
&self,
left: &Algebra,
right: &Algebra,
dataset: &dyn Dataset,
) -> Result<Solution> {
if let Some(result) = self.try_bound_join(left, right, dataset)? {
return Ok(result);
}
if let Some(result) = self.try_bound_join(right, left, dataset)? {
return Ok(result);
}
let left_results = self.execute_serial(left, dataset)?;
let right_results = self.execute_serial(right, dataset)?;
if left_results.len() < right_results.len() {
self.hash_join(left_results, right_results)
} else {
self.hash_join(right_results, left_results)
}
}
pub(super) fn try_bound_join(
&self,
outer: &Algebra,
inner: &Algebra,
dataset: &dyn Dataset,
) -> Result<Option<Solution>> {
if !is_bgp_or_filter_over_bgp(inner) {
return Ok(None);
}
let Some(outer_rows_hint) = static_row_count(outer) else {
return Ok(None);
};
if outer_rows_hint > BOUND_JOIN_ROW_THRESHOLD {
return Ok(None);
}
let pattern_vars = collect_pattern_variables(inner);
let outer_rows = self.execute_serial(outer, dataset)?;
let binds_pattern_var = outer_rows
.iter()
.any(|row| row.keys().any(|v| pattern_vars.contains(v)));
if !binds_pattern_var {
return Ok(None);
}
let Some(condition_vars) = collect_filter_condition_variables(inner) else {
return Ok(None);
};
let condition_only_var_bound = outer_rows.iter().any(|row| {
row.keys()
.any(|v| condition_vars.contains(v) && !pattern_vars.contains(v))
});
if condition_only_var_bound {
return Ok(None);
}
let mut result = Solution::new();
for outer_row in &outer_rows {
let specialized = substitute_algebra_binding(inner, outer_row);
let inner_rows = self.execute_serial(&specialized, dataset)?;
for inner_row in inner_rows {
if let Some(merged) = compatible_merge(outer_row, &inner_row) {
result.push(merged);
}
}
}
Ok(Some(result))
}
pub(super) fn try_filter_in_pushdown(
&self,
pattern: &Algebra,
condition: &crate::algebra::Expression,
dataset: &dyn Dataset,
) -> Result<Option<Solution>> {
use crate::algebra::{BinaryOperator, Expression};
let Expression::Binary {
op: BinaryOperator::In,
left,
right,
} = condition
else {
return Ok(None);
};
let Expression::Variable(v) = left.as_ref() else {
return Ok(None);
};
let Some(iris) = static_iri_list(right) else {
return Ok(None);
};
if iris.len() > BOUND_JOIN_ROW_THRESHOLD {
return Ok(None);
}
if !is_bgp_or_filter_over_bgp(pattern) {
return Ok(None);
}
if !collect_pattern_variables(pattern).contains(v) {
return Ok(None);
}
let mut seen: HashSet<Term> = HashSet::new();
let mut result = Solution::new();
for iri in iris {
if !seen.insert(iri.clone()) {
continue;
}
let mut row = Binding::new();
row.insert(v.clone(), iri);
let specialized = substitute_algebra_binding(pattern, &row);
for inner_row in self.execute_serial(&specialized, dataset)? {
if let Some(merged) = compatible_merge(&row, &inner_row) {
result.push(merged);
}
}
}
Ok(Some(result))
}
}
fn static_iri_list(expr: &crate::algebra::Expression) -> Option<Vec<Term>> {
use crate::algebra::Expression as E;
let items: &[E] = match expr {
E::Function { name, args }
if name.eq_ignore_ascii_case("list") || name.eq_ignore_ascii_case("in") =>
{
args
}
single => std::slice::from_ref(single),
};
let mut out = Vec::with_capacity(items.len());
for item in items {
match item {
E::Iri(iri) => out.push(Term::Iri(iri.clone())),
_ => return None,
}
}
Some(out)
}
fn is_bgp_or_filter_over_bgp(algebra: &Algebra) -> bool {
match algebra {
Algebra::Bgp(_) => true,
Algebra::Filter { pattern, .. } => is_bgp_or_filter_over_bgp(pattern),
_ => false,
}
}
fn static_row_count(algebra: &Algebra) -> Option<usize> {
match algebra {
Algebra::Values { bindings, .. } => Some(bindings.len()),
Algebra::Table => Some(1),
Algebra::Zero | Algebra::Empty => Some(0),
Algebra::Extend { pattern, .. } => static_row_count(pattern),
_ => None,
}
}
fn collect_pattern_variables(algebra: &Algebra) -> HashSet<Variable> {
let mut vars = HashSet::new();
collect_pattern_variables_into(algebra, &mut vars);
vars
}
fn collect_pattern_variables_into(algebra: &Algebra, vars: &mut HashSet<Variable>) {
match algebra {
Algebra::Bgp(triples) => {
for t in triples {
for term in [&t.subject, &t.predicate, &t.object] {
if let Term::Variable(v) = term {
vars.insert(v.clone());
}
}
}
}
Algebra::Filter { pattern, .. } => collect_pattern_variables_into(pattern, vars),
_ => {}
}
}
fn collect_filter_condition_variables(algebra: &Algebra) -> Option<HashSet<Variable>> {
let mut vars = HashSet::new();
let mut node = algebra;
while let Algebra::Filter { pattern, condition } = node {
if !collect_expression_variables_into(condition, &mut vars) {
return None;
}
node = pattern;
}
Some(vars)
}
fn collect_expression_variables_into(
expr: &crate::algebra::Expression,
vars: &mut HashSet<Variable>,
) -> bool {
use crate::algebra::Expression as E;
match expr {
E::Variable(v) | E::Bound(v) => {
vars.insert(v.clone());
true
}
E::Literal(_) | E::Iri(_) => true,
E::Function { args, .. } => args
.iter()
.all(|a| collect_expression_variables_into(a, vars)),
E::Binary { left, right, .. } => {
collect_expression_variables_into(left, vars)
&& collect_expression_variables_into(right, vars)
}
E::Unary { operand, .. } => collect_expression_variables_into(operand, vars),
E::Conditional {
condition,
then_expr,
else_expr,
} => {
collect_expression_variables_into(condition, vars)
&& collect_expression_variables_into(then_expr, vars)
&& collect_expression_variables_into(else_expr, vars)
}
E::Exists(_) | E::NotExists(_) => false,
}
}
fn compatible_merge(outer: &Binding, inner: &Binding) -> Option<Binding> {
let mut merged = outer.clone();
for (var, term) in inner {
if let Some(existing) = merged.get(var) {
if existing != term {
return None;
}
} else {
merged.insert(var.clone(), term.clone());
}
}
Some(merged)
}