use super::aggregate::aggregate;
use super::expr::SortKey;
use super::path::eval_path;
use super::*;
use crate::bgp::{
bgp_exists, collect_pattern_slots, eval_bgp_rows, row_to_binding, BgpSolutions, Binding,
PatternTerm, ProbeJoin, ProbePlan, TriplePattern,
};
use crate::file::Rete;
use crate::index::{GraphIndex, GraphIndexBuilder};
use crate::row::{bound_mask, merge_rows, Ctx, Row, Slots, Val};
use spargebra::term::{NamedNodePattern, TermPattern, TriplePattern as SpTriplePattern};
pub(super) type RowIter<'q> = Box<dyn Iterator<Item = Row> + 'q>;
const INLJ_MAX_HINT: usize = 4096;
fn inlj_hint(ctx: &Ctx) -> Option<usize> {
ctx.limit_hint.get().filter(|&h| h <= INLJ_MAX_HINT)
}
fn shares_certain_var(ctx: &Ctx, patterns: &[TriplePattern], lcert: &[bool]) -> bool {
patterns
.iter()
.flat_map(|p| [&p.s, &p.p, &p.o])
.any(|t| match t {
PatternTerm::Var(v) => ctx.slots.slot(v).is_some_and(|s| lcert[s]),
PatternTerm::Const(_) => false,
})
}
pub(super) fn query_ctx<'a>(rete: &'a Rete, sel: &Select) -> Ctx<'a> {
let mut slots = Slots::new();
collect_plan_slots(&sel.plan, &mut slots);
for (var, e) in &sel.extends {
collect_expr_slots(e, &mut slots);
slots.add(var);
}
if let Some(g) = &sel.group {
for (var, e) in &g.pre {
collect_expr_slots(e, &mut slots);
slots.add(var);
}
for v in &g.by {
slots.add(v);
}
for (res_var, agg) in &g.aggs {
slots.add(res_var);
match agg {
Agg::CountStar { .. } => {}
Agg::Count(v, _)
| Agg::Sum(v)
| Agg::Avg(v)
| Agg::Min(v)
| Agg::Max(v)
| Agg::Sample(v)
| Agg::GroupConcat(v, _, _) => {
slots.add(v);
}
}
}
}
for h in &sel.having {
collect_expr_slots(h, &mut slots);
}
for (e, _) in &sel.order {
collect_expr_slots(e, &mut slots);
}
for v in &sel.project {
slots.add(v);
}
Ctx::new(rete, slots)
}
fn collect_plan_slots(plan: &Plan, slots: &mut Slots) {
match plan {
Plan::Bgp(patterns) => collect_pattern_slots(patterns, slots),
Plan::Join(l, r) | Plan::Union(l, r) | Plan::Minus(l, r) => {
collect_plan_slots(l, slots);
collect_plan_slots(r, slots);
}
Plan::LeftJoin(l, r, cond) => {
collect_plan_slots(l, slots);
collect_plan_slots(r, slots);
if let Some(e) = cond {
collect_expr_slots(e, slots);
}
}
Plan::Filter(e, inner) => {
collect_expr_slots(e, slots);
collect_plan_slots(inner, slots);
}
Plan::Extend(var, e, inner) => {
slots.add(var);
collect_expr_slots(e, slots);
collect_plan_slots(inner, slots);
}
Plan::Subquery(sub) => {
if sub.project.is_empty() {
collect_plan_slots(&sub.plan, slots);
for (v, _) in &sub.extends {
slots.add(v);
}
if let Some(g) = &sub.group {
for v in &g.by {
slots.add(v);
}
for (rv, _) in &g.aggs {
slots.add(rv);
}
}
} else {
for v in &sub.project {
slots.add(v);
}
}
}
Plan::Path(s, _, o) => {
for t in [s, o] {
if let PatternTerm::Var(v) = t {
slots.add(v);
}
}
}
Plan::Values(vars, _) => {
for v in vars {
slots.add(v);
}
}
Plan::Service { vars, .. } => {
for v in vars {
slots.add(v);
}
}
Plan::Graph(target, inner) => {
if let GraphTarget::Var(v) = target {
slots.add(v);
}
collect_plan_slots(inner, slots);
}
}
}
fn collect_expr_slots(e: &FExpr, slots: &mut Slots) {
match e {
FExpr::Var(v) | FExpr::Bound(v) => {
slots.add(v);
}
FExpr::Const(_) => {}
FExpr::Arith(_, l, r)
| FExpr::Compare(_, l, r)
| FExpr::And(l, r)
| FExpr::Or(l, r)
| FExpr::SameTerm(l, r) => {
collect_expr_slots(l, slots);
collect_expr_slots(r, slots);
}
FExpr::Not(inner) => collect_expr_slots(inner, slots),
FExpr::If(c, t, e) => {
collect_expr_slots(c, slots);
collect_expr_slots(t, slots);
collect_expr_slots(e, slots);
}
FExpr::In(e, list) => {
collect_expr_slots(e, slots);
for x in list {
collect_expr_slots(x, slots);
}
}
FExpr::Func(_, args) | FExpr::Coalesce(args) => {
for a in args {
collect_expr_slots(a, slots);
}
}
FExpr::Exists(plan) => collect_plan_slots(plan, slots),
}
}
fn certain_bound(ctx: &Ctx, plan: &Plan, n: usize) -> Vec<bool> {
let mut m = vec![false; n];
mark_certain(ctx, plan, &mut m);
m
}
fn mark_certain(ctx: &Ctx, plan: &Plan, m: &mut [bool]) {
let mark_var = |v: &str, m: &mut [bool]| {
if let Some(i) = ctx.slots.slot(v) {
m[i] = true;
}
};
match plan {
Plan::Bgp(patterns) => {
for p in patterns {
for t in [&p.s, &p.p, &p.o] {
if let PatternTerm::Var(v) = t {
mark_var(v, m);
}
}
}
}
Plan::Path(s, _, o) => {
for t in [s, o] {
if let PatternTerm::Var(v) = t {
mark_var(v, m);
}
}
}
Plan::Values(vars, rows) => {
for (vi, v) in vars.iter().enumerate() {
if rows
.iter()
.all(|row| row.get(vi).is_some_and(Option::is_some))
{
mark_var(v, m);
}
}
}
Plan::Filter(_, inner) => mark_certain(ctx, inner, m),
Plan::Extend(_, _, inner) => mark_certain(ctx, inner, m),
Plan::Subquery(sub) => {
if sub.project.is_empty() {
mark_certain(ctx, &sub.plan, m);
} else {
for v in &sub.project {
mark_var(v, m);
}
}
for (v, _) in &sub.extends {
mark_var(v, m);
}
if let Some(g) = &sub.group {
for (rv, _) in &g.aggs {
mark_var(rv, m);
}
for v in &g.by {
mark_var(v, m);
}
}
}
Plan::Union(l, r) => {
let a = certain_bound(ctx, l, m.len());
let b = certain_bound(ctx, r, m.len());
for (i, slot) in m.iter_mut().enumerate() {
*slot |= a[i] && b[i];
}
}
Plan::Join(l, r) => {
mark_certain(ctx, l, m);
mark_certain(ctx, r, m);
}
Plan::LeftJoin(l, _, _) | Plan::Minus(l, _) => mark_certain(ctx, l, m),
Plan::Service { .. } => {}
Plan::Graph(target, inner) => {
if let GraphTarget::Var(v) = target {
mark_var(v, m);
}
mark_certain(ctx, inner, m);
}
}
}
fn possible_bound(ctx: &Ctx, plan: &Plan, n: usize) -> Vec<bool> {
let mut m = vec![false; n];
mark_possible(ctx, plan, &mut m);
m
}
fn mark_possible(ctx: &Ctx, plan: &Plan, m: &mut [bool]) {
let mark_var = |v: &str, m: &mut [bool]| {
if let Some(i) = ctx.slots.slot(v) {
m[i] = true;
}
};
match plan {
Plan::Bgp(patterns) => {
for p in patterns {
for t in [&p.s, &p.p, &p.o] {
if let PatternTerm::Var(v) = t {
mark_var(v, m);
}
}
}
}
Plan::Path(s, _, o) => {
for t in [s, o] {
if let PatternTerm::Var(v) = t {
mark_var(v, m);
}
}
}
Plan::Values(vars, _) => {
for v in vars {
mark_var(v, m);
}
}
Plan::Filter(_, inner) => mark_possible(ctx, inner, m),
Plan::Extend(var, _, inner) => {
mark_var(var, m);
mark_possible(ctx, inner, m);
}
Plan::Subquery(sub) => {
if sub.project.is_empty() {
mark_possible(ctx, &sub.plan, m);
} else {
for v in &sub.project {
mark_var(v, m);
}
}
for (v, _) in &sub.extends {
mark_var(v, m);
}
if let Some(g) = &sub.group {
for (rv, _) in &g.aggs {
mark_var(rv, m);
}
for v in &g.by {
mark_var(v, m);
}
}
}
Plan::Union(l, r) | Plan::Join(l, r) | Plan::LeftJoin(l, r, _) => {
mark_possible(ctx, l, m);
mark_possible(ctx, r, m);
}
Plan::Minus(l, _) => mark_possible(ctx, l, m),
Plan::Service { vars, .. } => {
for v in vars {
mark_var(v, m);
}
}
Plan::Graph(target, inner) => {
if let GraphTarget::Var(v) = target {
mark_var(v, m);
}
mark_possible(ctx, inner, m);
}
}
}
pub(super) fn ask_solution(rete: &Rete, sel: &Select) -> bool {
let ctx = query_ctx(rete, sel);
if sel.group.is_some() || !sel.having.is_empty() || !sel.extends.is_empty() {
return !raw_solutions_in(&ctx, sel).is_empty();
}
ctx.limit_hint.set(Some(1));
let mut merged = None;
let active = active_index(rete, &sel.from, &mut merged);
plan_exists(&ctx, active, sel.from_named.as_deref(), &sel.plan)
}
fn plan_exists(ctx: &Ctx, index: &GraphIndex, nf: Option<&[String]>, plan: &Plan) -> bool {
match plan {
Plan::Bgp(patterns) if patterns.len() <= 1 => bgp_exists(ctx, index, patterns),
Plan::Union(l, r) => plan_exists(ctx, index, nf, l) || plan_exists(ctx, index, nf, r),
Plan::Values(_, rows) => !rows.is_empty(),
_ => eval_plan_iter(ctx, index, nf, plan).next().is_some(),
}
}
pub(super) fn raw_solutions<'a>(rete: &'a Rete, sel: &Select) -> (Ctx<'a>, Vec<Row>) {
let ctx = query_ctx(rete, sel);
let rows = raw_solutions_in(&ctx, sel);
(ctx, rows)
}
fn raw_solutions_in(ctx: &Ctx, sel: &Select) -> Vec<Row> {
let mut merged = None;
let active = active_index(ctx.rete, &sel.from, &mut merged);
let nf = sel.from_named.as_deref();
let mut raw = match &sel.group {
Some(g) => aggregate(ctx, eval_plan_iter(ctx, active, nf, &sel.plan), g),
None => eval_plan_iter(ctx, active, nf, &sel.plan).collect(),
};
for row in raw.iter_mut() {
apply_extends_row(ctx, row, &sel.extends);
}
if !sel.having.is_empty() {
let mut cache = ExistsCache::new();
raw.retain(|b| {
sel.having
.iter()
.all(|f| f.boolean(ctx, active, b, &mut cache))
});
}
raw
}
fn apply_extends_row(ctx: &Ctx, row: &mut Row, extends: &[(String, FExpr)]) {
for (var, expr) in extends.iter().rev() {
if let Some(slot) = ctx.slots.slot(var) {
if let Some(v) = expr.value(ctx, row) {
row[slot] = Some(ctx.resolver.canon_term(&v));
}
}
}
}
fn active_index<'a>(
rete: &'a Rete,
from: &[String],
merged: &'a mut Option<GraphIndex>,
) -> &'a GraphIndex {
match from {
[] => rete.default_index(),
[g] => match rete.graph_index(g) {
Some(gi) => gi,
None => &*merged.insert(merge_graphs(rete, from)),
},
_ => &*merged.insert(merge_graphs(rete, from)),
}
}
fn merge_graphs(rete: &Rete, graphs: &[String]) -> GraphIndex {
let mut b = GraphIndexBuilder::new();
for g in graphs {
if let Some(gi) = rete.graph_index(g) {
for t in gi.match_pattern((None, None, None)) {
b.push(t);
}
}
}
b.build()
}
pub(super) fn instantiate(
ctx: &Ctx,
template: &[SpTriplePattern],
sols: &[Row],
) -> Vec<(String, String, String)> {
let mut set = std::collections::BTreeSet::new();
for b in sols {
for tp in template {
if let (Some(s), Some(p), Some(o)) = (
inst_term(ctx, &tp.subject, b),
inst_named(ctx, &tp.predicate, b),
inst_term(ctx, &tp.object, b),
) {
set.insert((s, p, o));
}
}
}
set.into_iter().collect()
}
fn row_var(ctx: &Ctx, name: &str, b: &Row) -> Option<String> {
let slot = ctx.slots.slot(name)?;
b[slot]
.as_ref()
.and_then(|v| ctx.resolver.str_of(v))
.map(|t| t.to_string())
}
fn inst_term(ctx: &Ctx, t: &TermPattern, b: &Row) -> Option<String> {
match t {
TermPattern::NamedNode(n) => Some(n.to_string()),
TermPattern::Literal(l) => Some(l.to_string()),
TermPattern::BlankNode(bn) => Some(bn.to_string()),
TermPattern::Variable(v) => row_var(ctx, v.as_str(), b),
TermPattern::Triple(tp) => {
let s = inst_term(ctx, &tp.subject, b)?;
let p = inst_named(ctx, &tp.predicate, b)?;
let o = inst_term(ctx, &tp.object, b)?;
Some(format!("<<{s} {p} {o}>>"))
}
}
}
fn inst_named(ctx: &Ctx, n: &NamedNodePattern, b: &Row) -> Option<String> {
match n {
NamedNodePattern::NamedNode(nn) => Some(nn.to_string()),
NamedNodePattern::Variable(v) => row_var(ctx, v.as_str(), b),
}
}
pub(super) fn run_select(rete: &Rete, sel: &Select) -> (Vec<String>, Vec<Binding>) {
let ctx = query_ctx(rete, sel);
if sel.order.is_empty() && !sel.distinct && sel.group.is_none() && sel.having.is_empty() {
ctx.limit_hint
.set(sel.limit.map(|l| l.saturating_add(sel.offset)));
}
let mut merged = None;
let active = active_index(rete, &sel.from, &mut merged);
let nf = sel.from_named.as_deref();
let source = eval_plan_iter(&ctx, active, nf, &sel.plan);
finish_select(&ctx, active, sel, source)
}
fn finish_select<'a, 'q>(
ctx: &'q Ctx<'a>,
active: &'q GraphIndex,
sel: &'q Select,
source: RowIter<'q>,
) -> (Vec<String>, Vec<Binding>) {
let mut source: RowIter<'q> = match &sel.group {
Some(g) => Box::new(aggregate(ctx, source, g).into_iter()),
None => source,
};
if !sel.extends.is_empty() {
let extends = &sel.extends;
source = Box::new(source.map(move |mut row| {
apply_extends_row(ctx, &mut row, extends);
row
}));
}
if !sel.having.is_empty() {
let having = &sel.having;
let mut cache = ExistsCache::new();
source = Box::new(
source.filter(move |b| having.iter().all(|f| f.boolean(ctx, active, b, &mut cache))),
);
}
if !sel.order.is_empty() {
let sorted = match (sel.limit, sel.distinct) {
(Some(limit), false) => {
top_k(ctx, source, &sel.order, sel.offset.saturating_add(limit))
}
_ => sort_all(ctx, source, &sel.order),
};
source = Box::new(sorted.into_iter());
}
let proj_slots: Vec<usize> = sel
.project
.iter()
.filter_map(|v| ctx.slots.slot(v))
.collect();
if !sel.project.is_empty() && sel.distinct {
let ps = proj_slots.clone();
source = Box::new(source.map(move |b| {
let mut p = ctx.slots.empty_row();
for &slot in &ps {
p[slot] = b[slot].clone();
}
p
}));
}
if sel.distinct {
let mut seen: std::collections::HashSet<Row> = std::collections::HashSet::new();
source = Box::new(source.filter(move |row| seen.insert(row.clone())));
}
let raw: Vec<Row> = source
.skip(sel.offset)
.take(sel.limit.unwrap_or(usize::MAX))
.collect();
if sel.project.is_empty() {
ctx.resolver
.prefetch(raw.iter().flat_map(|r| r.iter().filter_map(|v| v.as_ref())));
} else {
ctx.resolver.prefetch(
raw.iter()
.flat_map(|r| proj_slots.iter().filter_map(|&slot| r[slot].as_ref())),
);
}
let rows: Vec<Binding> = raw
.into_iter()
.map(|row| {
if sel.project.is_empty() {
row_to_binding(ctx, &row)
} else {
let mut b = Binding::new();
for (v, &slot) in sel.project.iter().zip(&proj_slots) {
if let Some(val) = &row[slot] {
if let Some(t) = ctx.resolver.str_once(val) {
b.insert(v.clone(), t);
}
}
}
b
}
})
.collect();
(sel.project.clone(), rows)
}
struct CommunityMembers {
community: usize,
subjects: usize,
members: Vec<Vec<Option<String>>>,
}
#[derive(Default)]
struct SplitStats {
rows_by_community: std::collections::BTreeMap<usize, usize>,
split_any: bool,
}
fn all_bound_mask(rows: &[Row], n: usize) -> Vec<bool> {
let mut m = vec![!rows.is_empty(); n];
for r in rows {
for (i, slot) in m.iter_mut().enumerate() {
*slot &= r[i].is_some();
}
}
m
}
fn join_rows(
ctx: &Ctx,
active: &GraphIndex,
left: Vec<Row>,
right: Vec<Row>,
optional: bool,
cond: Option<&FExpr>,
) -> Vec<Row> {
use std::collections::HashMap;
if right.is_empty() {
return if optional { left } else { Vec::new() };
}
let n = ctx.slots.len();
let lmask = all_bound_mask(&left, n);
let rmask = all_bound_mask(&right, n);
let key: Vec<usize> = (0..n).filter(|&i| lmask[i] && rmask[i]).collect();
let mut buckets: HashMap<Vec<Val>, Vec<usize>> = HashMap::new();
let mut partial: Vec<usize> = Vec::new();
for (i, row) in right.iter().enumerate() {
match key
.iter()
.map(|&s| row[s].clone())
.collect::<Option<Vec<Val>>>()
{
Some(k) => buckets.entry(k).or_default().push(i),
None => partial.push(i),
}
}
let mut cache = ExistsCache::new();
let mut out = Vec::new();
for lb in left {
let candidates: Vec<usize> = match key
.iter()
.map(|&s| lb[s].clone())
.collect::<Option<Vec<Val>>>()
{
Some(k) => buckets
.get(&k)
.into_iter()
.flatten()
.chain(partial.iter())
.copied()
.collect(),
None => (0..right.len()).collect(),
};
let mut matched = false;
for i in candidates {
if let Some(m) = merge_rows(&lb, &right[i]) {
if cond.is_none_or(|f| f.boolean(ctx, active, &m, &mut cache)) {
matched = true;
out.push(m);
}
}
}
if optional && !matched {
out.push(lb);
}
}
out
}
fn minus_rows(ctx: &Ctx, left: Vec<Row>, right: Vec<Row>) -> Vec<Row> {
use std::collections::HashMap;
if right.is_empty() {
return left;
}
let n = ctx.slots.len();
let lmask = all_bound_mask(&left, n);
let rmask = bound_mask(&right, n);
let key: Vec<usize> = (0..n).filter(|&i| lmask[i] && rmask[i]).collect();
let mut buckets: HashMap<Vec<Val>, Vec<usize>> = HashMap::new();
let mut partial: Vec<usize> = Vec::new();
for (i, row) in right.iter().enumerate() {
match key
.iter()
.map(|&s| row[s].clone())
.collect::<Option<Vec<Val>>>()
{
Some(k) => buckets.entry(k).or_default().push(i),
None => partial.push(i),
}
}
left.into_iter()
.filter(|lb| {
let eliminated = match key
.iter()
.map(|&s| lb[s].clone())
.collect::<Option<Vec<Val>>>()
{
Some(k) => {
buckets
.get(&k)
.is_some_and(|c| c.iter().any(|&i| minus_compatible(lb, &right[i])))
|| partial.iter().any(|&i| minus_compatible(lb, &right[i]))
}
None => right.iter().any(|rb| minus_compatible(lb, rb)),
};
!eliminated
})
.collect()
}
fn eval_split(
ctx: &Ctx,
active: &GraphIndex,
plan: &Plan,
parts: &[CommunityMembers],
stats: &mut SplitStats,
) -> Vec<Row> {
match plan {
Plan::Bgp(pats) => {
let mut groups: std::collections::BTreeMap<&str, Vec<TriplePattern>> =
std::collections::BTreeMap::new();
let mut residue: Vec<TriplePattern> = Vec::new();
for p in pats {
match &p.s {
PatternTerm::Var(v) => groups.entry(v.as_str()).or_default().push(p.clone()),
PatternTerm::Const(_) => residue.push(p.clone()),
}
}
if groups.is_empty() {
return eval_plan_in(ctx, active, None, plan);
}
stats.split_any = true;
let mut pieces: Vec<Vec<Row>> = Vec::new();
for (var, star) in &groups {
let star_plan = Plan::Bgp(star.clone());
let mut rows: Vec<Row> = Vec::new();
for part in parts {
let plan_c = Plan::Join(
Box::new(Plan::Values(vec![var.to_string()], part.members.clone())),
Box::new(star_plan.clone()),
);
let before = rows.len();
rows.extend(eval_plan_iter(ctx, active, None, &plan_c));
*stats.rows_by_community.entry(part.community).or_default() +=
rows.len() - before;
}
pieces.push(rows);
}
if !residue.is_empty() {
pieces.push(eval_plan_in(ctx, active, None, &Plan::Bgp(residue)));
}
pieces.sort_by_key(Vec::len);
let mut acc = pieces.remove(0);
for piece in pieces {
acc = join_rows(ctx, active, acc, piece, false, None);
}
acc
}
Plan::Filter(e, inner) => {
let rows = eval_split(ctx, active, inner, parts, stats);
let mut cache = ExistsCache::new();
rows.into_iter()
.filter(|b| e.boolean(ctx, active, b, &mut cache))
.collect()
}
Plan::Union(l, r) => {
let mut rows = eval_split(ctx, active, l, parts, stats);
rows.extend(eval_split(ctx, active, r, parts, stats));
rows
}
Plan::Join(l, r) => {
let lrows = eval_split(ctx, active, l, parts, stats);
let rrows = eval_split(ctx, active, r, parts, stats);
join_rows(ctx, active, lrows, rrows, false, None)
}
Plan::LeftJoin(l, r, cond) => {
let lrows = eval_split(ctx, active, l, parts, stats);
let rrows = eval_split(ctx, active, r, parts, stats);
join_rows(ctx, active, lrows, rrows, true, cond.as_ref())
}
Plan::Minus(l, r) => {
let lrows = eval_split(ctx, active, l, parts, stats);
let rrows = eval_split(ctx, active, r, parts, stats);
minus_rows(ctx, lrows, rrows)
}
_ => eval_plan_in(ctx, active, None, plan),
}
}
pub(super) fn run_select_communities(
rete: &Rete,
sel: &Select,
round: Option<usize>,
) -> Result<CommunitySelect, SparqlError> {
if !sel.from.is_empty() || sel.from_named.is_some() {
return Err(SparqlError::Unsupported(
"community-split evaluation works on the default graph only (no FROM / FROM NAMED)",
));
}
let dict = rete.dictionary();
let ids = rete.match_ids((None, None, None));
let g = crate::pyramid::project_graph(dict, &ids);
let dend = crate::pyramid::build_dendrogram(&g);
let round = round.unwrap_or_else(|| {
crate::tiling::choose_round_for_budget(dict, &ids, &dend, crate::file::DEFAULT_TILE_BUDGET)
});
let tiles = crate::tiling::tile_by_community(dict, &ids, &dend, round);
let parts: Vec<CommunityMembers> = tiles
.iter()
.map(|tile| {
let subjects: std::collections::BTreeSet<u32> =
tile.triples.iter().map(|&(s, _, _)| s).collect();
let members: Vec<Vec<Option<String>>> = subjects
.iter()
.filter_map(|&s| dict.subject_term(s))
.map(|t| vec![Some(t)])
.collect();
CommunityMembers {
community: tile.community,
subjects: members.len(),
members,
}
})
.collect();
let ctx = query_ctx(rete, sel);
let active = rete.default_index();
let mut stats = SplitStats::default();
let all = eval_split(&ctx, active, &sel.plan, &parts, &mut stats);
if !stats.split_any {
return Err(SparqlError::Unsupported(
"nothing to split: the query has no basic graph pattern with a variable subject — \
run it with the whole-index strategy",
));
}
let partials: Vec<CommunityPartial> = parts
.iter()
.map(|p| CommunityPartial {
community: p.community,
subjects: p.subjects,
rows: stats
.rows_by_community
.get(&p.community)
.copied()
.unwrap_or(0),
})
.collect();
let (vars, rows) = finish_select(&ctx, active, sel, Box::new(all.into_iter()));
Ok((vars, rows, partials))
}
fn cmp_keyed(
order: &[(FExpr, bool)],
a: &(Vec<SortKey>, usize, Row),
b: &(Vec<SortKey>, usize, Row),
) -> std::cmp::Ordering {
for (i, (_, desc)) in order.iter().enumerate() {
let ord = a.0[i].cmp(&b.0[i]);
let ord = if *desc { ord.reverse() } else { ord };
if ord != std::cmp::Ordering::Equal {
return ord;
}
}
a.1.cmp(&b.1)
}
fn sort_all(ctx: &Ctx, rows: RowIter, order: &[(FExpr, bool)]) -> Vec<Row> {
let mut keyed: Vec<(Vec<SortKey>, usize, Row)> = rows
.enumerate()
.map(|(seq, b)| {
let keys = order
.iter()
.map(|(e, _)| SortKey::of(e.value(ctx, &b)))
.collect();
(keys, seq, b)
})
.collect();
keyed.sort_by(|a, b| cmp_keyed(order, a, b));
keyed.into_iter().map(|(_, _, b)| b).collect()
}
fn top_k(ctx: &Ctx, rows: RowIter, order: &[(FExpr, bool)], k: usize) -> Vec<Row> {
if k == 0 {
return Vec::new();
}
let mut top: Vec<(Vec<SortKey>, usize, Row)> = Vec::with_capacity(k + 1);
for (seq, b) in rows.enumerate() {
let keys: Vec<SortKey> = order
.iter()
.map(|(e, _)| SortKey::of(e.value(ctx, &b)))
.collect();
let entry = (keys, seq, b);
if top.len() >= k && cmp_keyed(order, &entry, &top[k - 1]) != std::cmp::Ordering::Less {
continue;
}
let pos =
top.partition_point(|e| cmp_keyed(order, e, &entry) != std::cmp::Ordering::Greater);
top.insert(pos, entry);
top.truncate(k);
}
top.into_iter().map(|(_, _, b)| b).collect()
}
fn values_rows(ctx: &Ctx, vars: &[String], rows: &[Vec<Option<String>>]) -> Vec<Row> {
let slots: Vec<Option<usize>> = vars.iter().map(|v| ctx.slots.slot(v)).collect();
rows.iter()
.map(|row| {
let mut r = ctx.slots.empty_row();
for (slot, val) in slots.iter().zip(row.iter()) {
if let (Some(i), Some(t)) = (slot, val) {
r[*i] = Some(ctx.resolver.canon_term(t));
}
}
r
})
.collect()
}
pub(crate) fn eval_plan_in(
ctx: &Ctx,
index: &GraphIndex,
named_filter: Option<&[String]>,
plan: &Plan,
) -> Vec<Row> {
let saved = ctx.limit_hint.replace(None);
let rows = eval_plan_iter(ctx, index, named_filter, plan).collect();
ctx.limit_hint.set(saved);
rows
}
pub(crate) fn eval_plan_iter<'q>(
ctx: &'q Ctx<'q>,
index: &'q GraphIndex,
named_filter: Option<&'q [String]>,
plan: &'q Plan,
) -> RowIter<'q> {
let visible = move |name: &str| named_filter.is_none_or(|f| f.iter().any(|g| g == name));
match plan {
Plan::Bgp(patterns) => {
if inlj_hint(ctx).is_some() && patterns.len() >= 2 {
return match ProbeJoin::new(ctx, index, patterns) {
Some(pj) => Box::new(pj),
None => Box::new(std::iter::empty()),
};
}
Box::new(BgpSolutions::new(ctx, index, patterns))
}
Plan::Path(subj, spec, obj) => Box::new(eval_path(ctx, index, subj, spec, obj).into_iter()),
Plan::Values(vars, rows) => Box::new(values_rows(ctx, vars, rows).into_iter()),
Plan::Filter(expr, inner) => {
if let Some(rows) = text_contains_pushdown(ctx, index, named_filter, expr, inner) {
return Box::new(rows.into_iter());
}
let mut cache = ExistsCache::new();
Box::new(
eval_plan_iter(ctx, index, named_filter, inner)
.filter(move |b| expr.boolean(ctx, index, b, &mut cache)),
)
}
Plan::Subquery(sub) => {
let (_vars, bindings) = run_select(ctx.rete, sub);
Box::new(bindings_to_rows(ctx, bindings).into_iter())
}
Plan::Service {
silent,
endpoint,
query,
..
} => {
let result = match ctx.rete.service_client() {
Some(client) => client.query(endpoint, query),
None => Err(format!(
"{endpoint}: no SERVICE client attached to this file handle \
(the host must provide one; the CLI and browser clients do)"
)),
};
match result {
Ok(bindings) => Box::new(bindings_to_rows(ctx, bindings).into_iter()),
Err(e) if *silent => {
let _ = e;
Box::new(std::iter::once(ctx.slots.empty_row()))
}
Err(e) => {
ctx.rete.record_service_error(&e);
Box::new(std::iter::empty())
}
}
}
Plan::Extend(var, expr, inner) => {
let slot = ctx.slots.slot(var);
Box::new(
eval_plan_iter(ctx, index, named_filter, inner).map(move |mut row| {
if let Some(slot) = slot {
row[slot] = expr.value(ctx, &row).map(|v| ctx.resolver.canon_term(&v));
}
row
}),
)
}
Plan::Union(l, r) => Box::new(
eval_plan_iter(ctx, index, named_filter, l).chain(eval_plan_iter(
ctx,
index,
named_filter,
r,
)),
),
Plan::Minus(l, r) => minus_iter(ctx, index, named_filter, l, r),
Plan::Join(l, r) => {
if let Some(v) = values_pushdown(ctx, index, l, r) {
return Box::new(v.into_iter());
}
join_iter(ctx, index, named_filter, l, r, false, None)
}
Plan::LeftJoin(l, r, cond) => {
join_iter(ctx, index, named_filter, l, r, true, cond.as_ref())
}
Plan::Graph(GraphTarget::Named(iri), inner) => match ctx.rete.graph_index(iri) {
Some(gi) if visible(iri) => eval_plan_iter(ctx, gi, named_filter, inner),
_ => Box::new(std::iter::empty()),
},
Plan::Graph(GraphTarget::Var(var), inner) => {
let Some(slot) = ctx.slots.slot(var) else {
return Box::new(std::iter::empty());
};
Box::new(
ctx.rete
.named_graphs()
.iter()
.filter(move |(name, _)| visible(name))
.flat_map(move |(name, gi)| {
let gval = ctx.resolver.canon_term(name);
eval_plan_iter(ctx, gi, named_filter, inner).filter_map(move |mut sol| {
match &sol[slot] {
Some(existing) if *existing != gval => None,
_ => {
sol[slot] = Some(gval.clone());
Some(sol)
}
}
})
}),
)
}
}
}
fn bindings_to_rows(ctx: &Ctx, bindings: Vec<Binding>) -> Vec<Row> {
bindings
.into_iter()
.map(|b| {
let mut row = ctx.slots.empty_row();
for (var, term) in &b {
if let Some(slot) = ctx.slots.slot(var) {
row[slot] = Some(ctx.resolver.canon_term(term));
}
}
row
})
.collect()
}
fn substitute_patterns(
patterns: &[TriplePattern],
input: &[(String, String)],
) -> Vec<TriplePattern> {
let sub = |t: &PatternTerm| -> PatternTerm {
match t {
PatternTerm::Var(v) => match input.iter().find(|(k, _)| k == v) {
Some((_, val)) => PatternTerm::Const(val.clone()),
None => t.clone(),
},
PatternTerm::Const(_) => t.clone(),
}
};
patterns
.iter()
.map(|p| TriplePattern {
s: sub(&p.s),
p: sub(&p.p),
o: sub(&p.o),
})
.collect()
}
const TEXT_PUSHDOWN_MAX_CANDIDATES: usize = 8192;
fn contains_var(e: &FExpr) -> Option<&str> {
match e {
FExpr::Var(v) => Some(v),
FExpr::Func(Builtin::Str | Builtin::LCase | Builtin::UCase, args) if args.len() == 1 => {
contains_var(&args[0])
}
_ => None,
}
}
fn required_contains(e: &FExpr, out: &mut Vec<(String, String)>) {
match e {
FExpr::And(a, b) => {
required_contains(a, out);
required_contains(b, out);
}
FExpr::Func(Builtin::Contains, args) if args.len() == 2 => {
if let (Some(v), FExpr::Const(c)) = (contains_var(&args[0]), &args[1]) {
if let Some(needle) = crate::terms::literal_lexical(c) {
out.push((v.to_string(), needle));
}
}
}
_ => {}
}
}
fn text_contains_pushdown(
ctx: &Ctx,
index: &GraphIndex,
nf: Option<&[String]>,
expr: &FExpr,
inner: &Plan,
) -> Option<Vec<Row>> {
use std::collections::BTreeSet;
let ti = ctx.rete.text_index()?;
let mut needles: Vec<(String, String)> = Vec::new();
required_contains(expr, &mut needles);
let mut seed: Option<(String, BTreeSet<u32>)> = None;
for (cv, needle) in &needles {
let Some(sv) = contains_subject_var(inner, cv) else {
continue;
};
let pieces: Vec<String> = crate::text_index::tokenize(needle).collect();
if pieces.is_empty() {
continue; }
let mut cands: Option<BTreeSet<u32>> = None;
let mut usable = true;
for piece in &pieces {
let Some(subs) = ti.substring(piece) else {
usable = false; break;
};
let set: BTreeSet<u32> = subs.into_iter().collect();
cands = Some(match cands {
None => set,
Some(acc) => acc.intersection(&set).copied().collect(),
});
}
if !usable {
continue;
}
let cands = cands.unwrap_or_default();
match &mut seed {
None => seed = Some((sv, cands)),
Some((v, acc)) if *v == sv => *acc = acc.intersection(&cands).copied().collect(),
Some(_) => {}
}
}
let (sv, cands) = seed?;
if cands.len() > TEXT_PUSHDOWN_MAX_CANDIDATES {
return None;
}
let dict = ctx.rete.dictionary();
let ids: Vec<u32> = cands.iter().copied().collect();
dict.prefetch_subject_terms(&ids);
let rows: Vec<Vec<Option<String>>> = cands
.iter()
.filter_map(|&id| dict.subject_term(id))
.map(|t| vec![Some(t)])
.collect();
if let Some(target) = find_seed_bgp(inner, &sv) {
crate::bgp::prefetch_subject_probes(ctx, index, target, &sv, &ids);
}
let values = Plan::Values(vec![sv.clone()], rows);
let seeded = inject_contains_seed(inner, &sv, &values)?;
let id_rows: Vec<Row> = eval_plan_iter(ctx, index, nf, &seeded).collect();
let nodes: Vec<u32> = {
let mut set = BTreeSet::new();
for row in &id_rows {
for v in row.iter().flatten() {
if let Val::Id(x) = v {
if *x >= 0 {
set.insert(*x as u32);
}
}
}
}
set.into_iter().collect()
};
dict.prefetch_node_terms(&nodes);
let mut cache = ExistsCache::new();
let out: Vec<Row> = id_rows
.into_iter()
.filter(|b| expr.boolean(ctx, index, b, &mut cache))
.collect();
Some(out)
}
fn find_seed_bgp<'p>(plan: &'p Plan, sv: &str) -> Option<&'p [TriplePattern]> {
match plan {
Plan::Bgp(patterns) => patterns
.iter()
.any(|p| matches!(&p.s, PatternTerm::Var(s) if s == sv))
.then_some(patterns.as_slice()),
Plan::Join(l, r) => find_seed_bgp(l, sv).or_else(|| find_seed_bgp(r, sv)),
Plan::LeftJoin(l, _, _) => find_seed_bgp(l, sv),
Plan::Filter(_, inner) => find_seed_bgp(inner, sv),
_ => None,
}
}
fn contains_subject_var(plan: &Plan, cv: &str) -> Option<String> {
match plan {
Plan::Bgp(patterns) => patterns.iter().find_map(|p| match (&p.s, &p.o) {
(PatternTerm::Var(s), PatternTerm::Var(o)) if o == cv => Some(s.clone()),
_ => None,
}),
Plan::Join(l, r) => contains_subject_var(l, cv).or_else(|| contains_subject_var(r, cv)),
Plan::LeftJoin(l, _, _) => contains_subject_var(l, cv),
Plan::Filter(_, inner) => contains_subject_var(inner, cv),
_ => None,
}
}
fn inject_contains_seed(plan: &Plan, sv: &str, values: &Plan) -> Option<Plan> {
match plan {
Plan::Bgp(patterns) => {
if patterns
.iter()
.any(|p| matches!(&p.s, PatternTerm::Var(s) if s == sv))
{
Some(Plan::Join(Box::new(values.clone()), Box::new(plan.clone())))
} else {
None
}
}
Plan::Join(l, r) => {
if let Some(nl) = inject_contains_seed(l, sv, values) {
Some(Plan::Join(Box::new(nl), r.clone()))
} else {
inject_contains_seed(r, sv, values).map(|nr| Plan::Join(l.clone(), Box::new(nr)))
}
}
Plan::LeftJoin(l, r, c) => inject_contains_seed(l, sv, values)
.map(|nl| Plan::LeftJoin(Box::new(nl), r.clone(), c.clone())),
Plan::Filter(e, inner) => {
inject_contains_seed(inner, sv, values).map(|ni| Plan::Filter(e.clone(), Box::new(ni)))
}
_ => None,
}
}
fn values_pushdown(ctx: &Ctx, index: &GraphIndex, l: &Plan, r: &Plan) -> Option<Vec<Row>> {
let (vals, patterns) = match (l, r) {
(Plan::Values(v, rows), Plan::Bgp(p)) | (Plan::Bgp(p), Plan::Values(v, rows)) => {
((v, rows), p)
}
_ => return None,
};
let (vars, rows) = vals;
let bgp_vars: std::collections::HashSet<&str> = patterns
.iter()
.flat_map(|p| [&p.s, &p.p, &p.o])
.filter_map(|t| match t {
PatternTerm::Var(v) => Some(v.as_str()),
PatternTerm::Const(_) => None,
})
.collect();
if !vars.iter().any(|v| bgp_vars.contains(v.as_str())) {
return None;
}
let mut out = Vec::new();
for row in rows {
let input: Vec<(String, String)> = vars
.iter()
.zip(row.iter())
.filter_map(|(v, val)| val.as_ref().map(|t| (v.clone(), t.clone())))
.collect();
let subst = substitute_patterns(patterns, &input);
let mut base = ctx.slots.empty_row();
for (v, t) in &input {
if let Some(i) = ctx.slots.slot(v) {
base[i] = Some(ctx.resolver.canon_term(t));
}
}
for brow in eval_bgp_rows(ctx, index, &subst) {
let mut merged = base.clone();
for (slot, v) in brow.iter().enumerate() {
if v.is_some() {
merged[slot] = v.clone();
}
}
out.push(merged);
}
}
Some(out)
}
fn minus_compatible(lb: &Row, rb: &Row) -> bool {
let mut shared = false;
for (l, r) in lb.iter().zip(rb.iter()) {
if let (Some(v), Some(w)) = (l, r) {
if v != w {
return false;
}
shared = true;
}
}
shared
}
fn minus_iter<'q>(
ctx: &'q Ctx<'q>,
index: &'q GraphIndex,
nf: Option<&'q [String]>,
l: &'q Plan,
r: &'q Plan,
) -> RowIter<'q> {
use std::collections::HashMap;
let right: Vec<Row> = eval_plan_iter(ctx, index, nf, r).collect();
if right.is_empty() {
return eval_plan_iter(ctx, index, nf, l);
}
let n = ctx.slots.len();
let rmask = bound_mask(&right, n);
let lposs = possible_bound(ctx, l, n);
if !(0..n).any(|i| lposs[i] && rmask[i]) {
return eval_plan_iter(ctx, index, nf, l);
}
let lcert = certain_bound(ctx, l, n);
let jv: Vec<usize> = (0..n).filter(|&i| lcert[i] && rmask[i]).collect();
let mut buckets: HashMap<Vec<Val>, Vec<usize>> = HashMap::new();
let mut partial: Vec<usize> = Vec::new();
for (i, row) in right.iter().enumerate() {
match jv
.iter()
.map(|&s| row[s].clone())
.collect::<Option<Vec<Val>>>()
{
Some(k) => buckets.entry(k).or_default().push(i),
None => partial.push(i),
}
}
let left = eval_plan_iter(ctx, index, nf, l);
Box::new(left.filter(move |lb| {
let eliminated = match jv
.iter()
.map(|&s| lb[s].clone())
.collect::<Option<Vec<Val>>>()
{
Some(k) => {
let in_bucket = buckets
.get(&k)
.is_some_and(|c| c.iter().any(|&i| minus_compatible(lb, &right[i])));
in_bucket || partial.iter().any(|&i| minus_compatible(lb, &right[i]))
}
None => right.iter().any(|rb| minus_compatible(lb, rb)),
};
!eliminated
}))
}
struct JoinIter<'q> {
ctx: &'q Ctx<'q>,
index: &'q GraphIndex,
left: RowIter<'q>,
right: Vec<Row>,
buckets: std::collections::HashMap<Vec<Val>, Vec<usize>>,
partial: Vec<usize>,
jv: Vec<usize>,
optional: bool,
cond: Option<&'q FExpr>,
cache: ExistsCache,
cur_left: Option<Row>,
candidates: Vec<usize>,
ci: usize,
matched: bool,
}
fn path_endpoint_bound(ctx: &Ctx, other: &Plan, s: &PatternTerm, o: &PatternTerm) -> bool {
let cb = certain_bound(ctx, other, ctx.slots.len());
[s, o].into_iter().any(|t| match t {
PatternTerm::Var(v) => ctx.slots.slot(v).is_some_and(|i| cb[i]),
PatternTerm::Const(_) => false,
})
}
fn fix_endpoint(ctx: &Ctx, t: &PatternTerm, lrow: &Row) -> PatternTerm {
if let PatternTerm::Var(v) = t {
if let Some(val) = ctx.slots.slot(v).and_then(|i| lrow[i].as_ref()) {
if let Some(term) = ctx.resolver.str_once(val) {
return PatternTerm::Const(term);
}
}
}
t.clone()
}
#[allow(clippy::too_many_arguments)]
fn correlated_path_join<'q>(
ctx: &'q Ctx<'q>,
index: &'q GraphIndex,
nf: Option<&'q [String]>,
bound: &'q Plan,
subj: &'q PatternTerm,
spec: &'q PathAst,
obj: &'q PatternTerm,
optional: bool,
cond: Option<&'q FExpr>,
) -> RowIter<'q> {
Box::new(eval_plan_iter(ctx, index, nf, bound).flat_map(move |lrow| {
let (s2, o2) = (
fix_endpoint(ctx, subj, &lrow),
fix_endpoint(ctx, obj, &lrow),
);
let mut cache = ExistsCache::new();
let mut out: Vec<Row> = Vec::new();
for pr in eval_path(ctx, index, &s2, spec, &o2) {
if let Some(m) = merge_rows(&lrow, &pr) {
if cond.is_none_or(|f| f.boolean(ctx, index, &m, &mut cache)) {
out.push(m);
}
}
}
if optional && out.is_empty() {
out.push(lrow);
}
out.into_iter()
}))
}
fn join_iter<'q>(
ctx: &'q Ctx<'q>,
index: &'q GraphIndex,
nf: Option<&'q [String]>,
l: &'q Plan,
r: &'q Plan,
optional: bool,
cond: Option<&'q FExpr>,
) -> RowIter<'q> {
if let Plan::Bgp(patterns) = r {
if !patterns.is_empty() {
let lcert = certain_bound(ctx, l, ctx.slots.len());
let probe = inlj_hint(ctx).is_some()
|| (shares_certain_var(ctx, patterns, &lcert)
&& crate::bgp::bgp_min_scan_bytes(ctx, index, patterns)
.is_none_or(|n| n >= crate::bgp::FAT_SCAN_BYTES));
if probe {
return match ProbePlan::new(ctx, patterns, &lcert) {
Some(plan) => {
let left: RowIter<'q> = if inlj_hint(ctx).is_none() {
let rows: Vec<Row> = eval_plan_iter(ctx, index, nf, l).collect();
crate::bgp::prefetch_plan_probes(ctx, index, &plan, &rows);
Box::new(rows.into_iter())
} else {
eval_plan_iter(ctx, index, nf, l)
};
Box::new(ProbedJoin {
ctx,
index,
left,
plan,
optional,
cond,
cache: ExistsCache::new(),
cur: None,
})
}
None if optional => eval_plan_iter(ctx, index, nf, l),
None => Box::new(std::iter::empty()),
};
}
}
}
if let Plan::Path(s, spec, o) = r {
if path_endpoint_bound(ctx, l, s, o) {
return correlated_path_join(ctx, index, nf, l, s, spec, o, optional, cond);
}
}
if !optional {
if let Plan::Path(s, spec, o) = l {
if path_endpoint_bound(ctx, r, s, o) {
return correlated_path_join(ctx, index, nf, r, s, spec, o, false, cond);
}
}
}
let right: Vec<Row> = eval_plan_iter(ctx, index, nf, r).collect();
if right.is_empty() {
return if optional {
eval_plan_iter(ctx, index, nf, l)
} else {
Box::new(std::iter::empty())
};
}
let n = ctx.slots.len();
let jv: Vec<usize> = {
let lcert = certain_bound(ctx, l, n);
let rcert = certain_bound(ctx, r, n);
let rmask = bound_mask(&right, n);
(0..n)
.filter(|&i| lcert[i] && rcert[i] && rmask[i])
.collect()
};
let mut buckets: std::collections::HashMap<Vec<Val>, Vec<usize>> =
std::collections::HashMap::new();
let mut partial: Vec<usize> = Vec::new();
for (i, row) in right.iter().enumerate() {
match jv
.iter()
.map(|&s| row[s].clone())
.collect::<Option<Vec<Val>>>()
{
Some(k) => buckets.entry(k).or_default().push(i),
None => partial.push(i),
}
}
Box::new(JoinIter {
ctx,
index,
left: eval_plan_iter(ctx, index, nf, l),
right,
buckets,
partial,
jv,
optional,
cond,
cache: ExistsCache::new(),
cur_left: None,
candidates: Vec::new(),
ci: 0,
matched: false,
})
}
struct ProbedJoin<'q> {
ctx: &'q Ctx<'q>,
index: &'q GraphIndex,
left: RowIter<'q>,
plan: ProbePlan,
optional: bool,
cond: Option<&'q FExpr>,
cache: ExistsCache,
cur: Option<(Row, ProbeJoin<'q>, bool)>,
}
impl Iterator for ProbedJoin<'_> {
type Item = Row;
fn next(&mut self) -> Option<Row> {
loop {
if let Some((_, probe, matched)) = &mut self.cur {
for m in probe.by_ref() {
if self
.cond
.is_none_or(|f| f.boolean(self.ctx, self.index, &m, &mut self.cache))
{
*matched = true;
return Some(m);
}
}
let (l, _, matched) = self.cur.take().unwrap();
if self.optional && !matched {
return Some(l);
}
}
let l = self.left.next()?;
let probe = ProbeJoin::from_plan(self.ctx, self.index, &self.plan, l.clone());
self.cur = Some((l, probe, false));
}
}
}
impl Iterator for JoinIter<'_> {
type Item = Row;
fn next(&mut self) -> Option<Row> {
loop {
if let Some(left) = &self.cur_left {
while self.ci < self.candidates.len() {
let ri = self.candidates[self.ci];
self.ci += 1;
if let Some(m) = merge_rows(left, &self.right[ri]) {
if self
.cond
.is_none_or(|f| f.boolean(self.ctx, self.index, &m, &mut self.cache))
{
self.matched = true;
return Some(m);
}
}
}
let l = self.cur_left.take().unwrap();
if self.optional && !self.matched {
return Some(l);
}
}
let l = self.left.next()?;
self.candidates = match self
.jv
.iter()
.map(|&s| l[s].clone())
.collect::<Option<Vec<Val>>>()
{
Some(k) => {
let mut c = self.buckets.get(&k).cloned().unwrap_or_default();
c.extend_from_slice(&self.partial);
c
}
None => (0..self.right.len()).collect(),
};
self.ci = 0;
self.matched = false;
self.cur_left = Some(l);
}
}
}