Skip to main content

oxirs_core/query/
mod.rs

1//! SPARQL query processing module with full SciRS2 integration
2
3pub mod advanced_statistics;
4pub mod algebra;
5pub mod binding_optimizer;
6pub mod cost_based_optimizer;
7pub mod cost_estimator;
8pub mod distributed;
9pub mod exec;
10pub mod functions;
11pub mod gpu;
12pub mod jit;
13pub mod ml_optimizer;
14pub mod optimizer;
15pub mod parser;
16pub mod pattern_optimizer;
17pub mod pattern_unification;
18pub mod plan;
19pub mod plan_cache;
20pub mod profiled_plan_builder;
21pub mod property_function_registry;
22pub mod property_paths;
23pub mod query_plan_visualizer;
24pub mod query_profiler;
25pub mod result_cache;
26pub mod sparql_algebra;
27pub mod sparql_algebra_ops;
28#[cfg(test)]
29mod sparql_algebra_tests;
30pub mod sparql_algebra_transform;
31pub mod sparql_algebra_types;
32pub mod sparql_algebra_types_expr;
33pub mod sparql_algebra_types_paths;
34pub mod sparql_algebra_types_pattern;
35pub mod sparql_algebra_types_terms;
36pub mod sparql_query;
37pub mod statistics;
38pub mod streaming_results;
39pub mod update;
40pub mod wasm;
41
42// Re-export property function registry types
43pub use property_function_registry::{
44    PropertyFunction, PropertyFunctionArg, PropertyFunctionBinding, PropertyFunctionFactory,
45    PropertyFunctionMetadata, PropertyFunctionRegistry, PropertyFunctionResult,
46};
47
48// Re-export the enhanced SPARQL algebra and query types from sparql_algebra
49pub use crate::{GraphName, Triple};
50pub use sparql_algebra::{
51    Expression as SparqlExpression, GraphPattern as SparqlGraphPattern, NamedNodePattern,
52    PropertyPathExpression, TermPattern as SparqlTermPattern, TriplePattern as SparqlTriplePattern,
53};
54pub use sparql_query::*;
55
56// Re-export execution plan types
57pub use plan::ExecutionPlan;
58
59// Re-export advanced statistics types
60pub use advanced_statistics::{AdvancedStatistics, AdvancedStatisticsCollector, PatternExecution};
61
62// Re-export algebra types (without Query to avoid conflict)
63// Use explicit aliases to avoid conflicts
64pub use algebra::{
65    AlgebraTriplePattern, Expression as AlgebraExpression, GraphPattern as AlgebraGraphPattern,
66    PropertyPath, Query as AlgebraQuery, TermPattern as AlgebraTermPattern,
67};
68pub use binding_optimizer::{BindingIterator, BindingOptimizer, BindingSet, Constraint, TermType};
69pub use cost_based_optimizer::{
70    CostBasedOptimizer, CostConfiguration, Optimization, OptimizedPlan, OptimizerStats,
71};
72pub use distributed::{DistributedConfig, DistributedQueryEngine, FederatedEndpoint};
73pub use gpu::{GpuBackend, GpuQueryExecutor};
74pub use jit::{JitCompiler, JitConfig};
75pub use ml_optimizer::{
76    MLOptimizationResult, MLOptimizerConfig, MLQueryOptimizer, PatternFeatures, PerformanceMetrics,
77    TrainingStats,
78};
79pub use optimizer::{AIQueryOptimizer, MultiQueryOptimizer};
80pub use parser::*;
81pub use pattern_optimizer::{IndexType, OptimizedPatternPlan, PatternExecutor, PatternOptimizer};
82pub use pattern_unification::{
83    PatternConverter, PatternOptimizer as UnifiedPatternOptimizer, UnifiedTermPattern,
84    UnifiedTriplePattern,
85};
86pub use plan_cache::{
87    CacheConfig, CacheStatistics, CachedPlan, LruQueryPlanCache, PlanCacheStats, QueryPlan,
88    QueryPlanCache, SerializablePlan,
89};
90pub use profiled_plan_builder::{
91    CacheEffectiveness, ExecutionComparison, ImprovementLevel, PerformanceAnalysis,
92    PerformanceGrade, ProfiledPlanBuilder, ProfilingReport,
93};
94pub use query_plan_visualizer::{
95    HintSeverity, OptimizationHint, QueryPlanNode, QueryPlanSummary, QueryPlanVisualizer,
96};
97pub use query_profiler::{
98    ProfiledQuery, ProfilerConfig, ProfilingStatistics, QueryProfiler, QueryProfilingSession,
99    QueryStatistics,
100};
101pub use result_cache::{CacheConfig as ResultCacheConfig, CacheStats, QueryResultCache};
102pub use statistics::{
103    GraphStatistics, PredicateStatistics, QueryExecutionStats, SelectivityInfo, StatisticsSummary,
104};
105pub use streaming_results::{
106    ConstructResults, SelectResults, Solution as StreamingSolution, SolutionMetadata,
107    StreamingConfig, StreamingProgress, StreamingQueryResults, StreamingResultBuilder,
108};
109pub use update::{UpdateExecutor, UpdateParser};
110pub use wasm::{OptimizationLevel, WasmQueryCompiler, WasmTarget};
111
112// TODO: Temporary compatibility layer for SHACL module
113pub use exec::{QueryExecutor, QueryResults, Solution};
114
115use crate::model::{Object, Predicate, Subject, Term, Variable};
116use crate::OxirsError;
117use crate::Store;
118use std::collections::HashMap;
119use std::future::Future;
120use std::pin::Pin;
121
122// Import TermPattern for internal usage
123use algebra::TermPattern;
124
125/// Type alias for federated query execution future
126type FederatedQueryFuture<'a> =
127    Pin<Box<dyn Future<Output = Result<Vec<HashMap<String, Term>>, OxirsError>> + 'a>>;
128
129/// Simplified QueryResult for SHACL compatibility
130#[derive(Debug, Clone)]
131pub enum QueryResult {
132    /// SELECT query results
133    Select {
134        variables: Vec<String>,
135        bindings: Vec<HashMap<String, Term>>,
136    },
137    /// ASK query results
138    Ask(bool),
139    /// CONSTRUCT query results
140    Construct(Vec<crate::model::Triple>),
141}
142
143/// Simplified QueryEngine for SHACL compatibility
144pub struct QueryEngine {
145    /// Query parser for converting SPARQL strings to Query objects
146    parser: parser::SparqlParser,
147    /// Query executor for executing plans
148    executor_config: QueryExecutorConfig,
149    /// Federation executor for SERVICE clause support
150    federation_executor: Option<crate::federation::FederationExecutor>,
151}
152
153impl Default for QueryEngine {
154    fn default() -> Self {
155        Self::new()
156    }
157}
158
159/// Configuration for query execution
160#[derive(Debug, Clone)]
161pub struct QueryExecutorConfig {
162    /// Maximum number of results to return
163    pub max_results: usize,
164    /// Query timeout in milliseconds
165    pub timeout_ms: Option<u64>,
166    /// Enable query optimization
167    pub optimize: bool,
168}
169
170impl Default for QueryExecutorConfig {
171    fn default() -> Self {
172        Self {
173            max_results: 10000,
174            timeout_ms: Some(30000),
175            optimize: true,
176        }
177    }
178}
179
180impl QueryEngine {
181    /// Create a new query engine
182    pub fn new() -> Self {
183        Self {
184            parser: parser::SparqlParser::new(),
185            executor_config: QueryExecutorConfig::default(),
186            federation_executor: crate::federation::FederationExecutor::new().ok(),
187        }
188    }
189
190    /// Create a new query engine with custom configuration
191    pub fn with_config(config: QueryExecutorConfig) -> Self {
192        Self {
193            parser: parser::SparqlParser::new(),
194            executor_config: config,
195            federation_executor: crate::federation::FederationExecutor::new().ok(),
196        }
197    }
198
199    /// Enable federation support
200    pub fn with_federation(mut self) -> Self {
201        self.federation_executor = crate::federation::FederationExecutor::new().ok();
202        self
203    }
204
205    /// Disable federation support
206    pub fn without_federation(mut self) -> Self {
207        self.federation_executor = None;
208        self
209    }
210
211    /// Execute a SPARQL query string against a store
212    pub fn query(&self, query_str: &str, store: &dyn Store) -> Result<QueryResult, OxirsError> {
213        // Parse the query string
214        let parsed_query = self.parser.parse(query_str)?;
215
216        // Execute the parsed query
217        self.execute_query(&parsed_query, store)
218    }
219
220    /// Execute a parsed Query object against a store
221    pub fn execute_query(
222        &self,
223        query: &sparql_query::Query,
224        store: &dyn Store,
225    ) -> Result<QueryResult, OxirsError> {
226        match query {
227            sparql_query::Query::Select {
228                pattern, dataset, ..
229            } => self.execute_select_query(pattern, dataset.as_ref(), store),
230            sparql_query::Query::Ask {
231                pattern, dataset, ..
232            } => self.execute_ask_query(pattern, dataset.as_ref(), store),
233            sparql_query::Query::Construct {
234                template,
235                pattern,
236                dataset,
237                ..
238            } => self.execute_construct_query(template, pattern, dataset.as_ref(), store),
239            sparql_query::Query::Describe {
240                pattern, dataset, ..
241            } => self.execute_describe_query(pattern, dataset.as_ref(), store),
242        }
243    }
244
245    /// Execute a SELECT query
246    fn execute_select_query(
247        &self,
248        pattern: &SparqlGraphPattern,
249        _dataset: Option<&QueryDataset>,
250        store: &dyn Store,
251    ) -> Result<QueryResult, OxirsError> {
252        let executor = QueryExecutor::new(store);
253
254        // Convert graph pattern to execution plan
255        let plan = self.pattern_to_plan(pattern)?;
256
257        // Execute the plan
258        let solutions = executor.execute(&plan)?;
259
260        // Extract the *result* variable names (honoring any explicit projection)
261        // and convert solutions.
262        let variables = self.result_variables(pattern);
263        let bindings: Vec<HashMap<String, Term>> = solutions
264            .into_iter()
265            .take(self.executor_config.max_results)
266            .map(|sol| {
267                let mut binding = HashMap::new();
268                for var in &variables {
269                    if let Some(term) = sol.get(var) {
270                        binding.insert(var.name().to_string(), term.clone());
271                    }
272                }
273                binding
274            })
275            .collect();
276
277        Ok(QueryResult::Select {
278            variables: variables
279                .into_iter()
280                .map(|v| v.name().to_string())
281                .collect(),
282            bindings,
283        })
284    }
285
286    /// Execute an ASK query
287    fn execute_ask_query(
288        &self,
289        pattern: &SparqlGraphPattern,
290        _dataset: Option<&QueryDataset>,
291        store: &dyn Store,
292    ) -> Result<QueryResult, OxirsError> {
293        let executor = QueryExecutor::new(store);
294
295        // Convert graph pattern to execution plan
296        let plan = self.pattern_to_plan(pattern)?;
297
298        // Execute the plan
299        let solutions = executor.execute(&plan)?;
300
301        // ASK query returns true if there are any solutions
302        Ok(QueryResult::Ask(!solutions.is_empty()))
303    }
304
305    /// Execute a CONSTRUCT query
306    fn execute_construct_query(
307        &self,
308        template: &[SparqlTriplePattern],
309        pattern: &SparqlGraphPattern,
310        _dataset: Option<&QueryDataset>,
311        store: &dyn Store,
312    ) -> Result<QueryResult, OxirsError> {
313        let executor = QueryExecutor::new(store);
314
315        // Convert graph pattern to execution plan
316        let plan = self.pattern_to_plan(pattern)?;
317
318        // Execute the plan
319        let solutions = executor.execute(&plan)?;
320
321        // Construct triples from template and solutions
322        let mut triples = Vec::new();
323        for solution in solutions.into_iter().take(self.executor_config.max_results) {
324            for triple_pattern in template {
325                if let Some(triple) = self.instantiate_triple_pattern(triple_pattern, &solution)? {
326                    triples.push(triple);
327                }
328            }
329        }
330
331        Ok(QueryResult::Construct(triples))
332    }
333
334    /// Execute a DESCRIBE query
335    fn execute_describe_query(
336        &self,
337        pattern: &SparqlGraphPattern,
338        _dataset: Option<&QueryDataset>,
339        store: &dyn Store,
340    ) -> Result<QueryResult, OxirsError> {
341        // For now, treat DESCRIBE like CONSTRUCT *
342        // This is a simplified implementation
343        let executor = QueryExecutor::new(store);
344
345        // Convert graph pattern to execution plan
346        let plan = self.pattern_to_plan(pattern)?;
347
348        // Execute the plan
349        let solutions = executor.execute(&plan)?;
350
351        // Get all triples involving the found entities
352        let mut triples = Vec::new();
353        for solution in solutions.into_iter().take(self.executor_config.max_results) {
354            // For each bound entity, get all triples where it appears
355            for (_, term) in solution.iter() {
356                if let Ok(store_quads) =
357                    store.find_quads(None, None, None, Some(&GraphName::DefaultGraph))
358                {
359                    for quad in store_quads {
360                        let triple = Triple::new(
361                            quad.subject().clone(),
362                            quad.predicate().clone(),
363                            quad.object().clone(),
364                        );
365                        if self.triple_involves_term(&triple, term) {
366                            triples.push(triple);
367                        }
368                    }
369                }
370            }
371        }
372
373        triples.dedup();
374        Ok(QueryResult::Construct(triples))
375    }
376
377    /// Convert a graph pattern to an execution plan
378    fn pattern_to_plan(&self, pattern: &SparqlGraphPattern) -> Result<ExecutionPlan, OxirsError> {
379        match pattern {
380            SparqlGraphPattern::Bgp { patterns } => {
381                if patterns.len() == 1 {
382                    // Single triple pattern
383                    Ok(ExecutionPlan::TripleScan {
384                        pattern: self.convert_sparql_triple_pattern(&patterns[0])?,
385                    })
386                } else {
387                    // Multiple patterns - join them
388                    let mut plan = ExecutionPlan::TripleScan {
389                        pattern: self.convert_sparql_triple_pattern(&patterns[0])?,
390                    };
391
392                    for triple_pattern in &patterns[1..] {
393                        let right_plan = ExecutionPlan::TripleScan {
394                            pattern: self.convert_sparql_triple_pattern(triple_pattern)?,
395                        };
396
397                        // Find join variables
398                        let join_vars = self.find_join_variables(&plan, &right_plan);
399
400                        plan = ExecutionPlan::HashJoin {
401                            left: Box::new(plan),
402                            right: Box::new(right_plan),
403                            join_vars,
404                        };
405                    }
406
407                    Ok(plan)
408                }
409            }
410            SparqlGraphPattern::Join { left, right } => {
411                let left_plan = self.pattern_to_plan(left)?;
412                let right_plan = self.pattern_to_plan(right)?;
413                let join_vars = self.find_join_variables(&left_plan, &right_plan);
414
415                Ok(ExecutionPlan::HashJoin {
416                    left: Box::new(left_plan),
417                    right: Box::new(right_plan),
418                    join_vars,
419                })
420            }
421            SparqlGraphPattern::Filter { expr, inner } => {
422                let input_plan = self.pattern_to_plan(inner)?;
423                // Convert sparql_algebra::Expression to algebra::Expression
424                let condition = self.convert_expression(expr.clone())?;
425                Ok(ExecutionPlan::Filter {
426                    input: Box::new(input_plan),
427                    condition,
428                })
429            }
430            SparqlGraphPattern::Union { left, right } => {
431                let left_plan = self.pattern_to_plan(left)?;
432                let right_plan = self.pattern_to_plan(right)?;
433
434                Ok(ExecutionPlan::Union {
435                    left: Box::new(left_plan),
436                    right: Box::new(right_plan),
437                })
438            }
439            SparqlGraphPattern::Project { inner, variables } => {
440                let input_plan = self.pattern_to_plan(inner)?;
441                Ok(ExecutionPlan::Project {
442                    input: Box::new(input_plan),
443                    vars: variables.clone(),
444                })
445            }
446            SparqlGraphPattern::Distinct { inner } => {
447                let input_plan = self.pattern_to_plan(inner)?;
448                Ok(ExecutionPlan::Distinct {
449                    input: Box::new(input_plan),
450                })
451            }
452            SparqlGraphPattern::Reduced { inner } => {
453                // REDUCED permits (but does not require) duplicate elimination.
454                // Implementing it as DISTINCT is a conformant choice.
455                let input_plan = self.pattern_to_plan(inner)?;
456                Ok(ExecutionPlan::Distinct {
457                    input: Box::new(input_plan),
458                })
459            }
460            SparqlGraphPattern::OrderBy { inner, expression } => {
461                let input_plan = self.pattern_to_plan(inner)?;
462                let order_by = expression
463                    .iter()
464                    .map(|oe| self.convert_order_expression(oe))
465                    .collect::<Result<Vec<_>, _>>()?;
466                Ok(ExecutionPlan::Sort {
467                    input: Box::new(input_plan),
468                    order_by,
469                })
470            }
471            SparqlGraphPattern::Slice {
472                inner,
473                start,
474                length,
475            } => {
476                let input_plan = self.pattern_to_plan(inner)?;
477                Ok(ExecutionPlan::Limit {
478                    input: Box::new(input_plan),
479                    limit: length.unwrap_or(usize::MAX),
480                    offset: *start,
481                })
482            }
483            _ => {
484                // For unsupported patterns, return an error for now
485                Err(OxirsError::Query(format!(
486                    "Unsupported graph pattern type: {pattern:?}"
487                )))
488            }
489        }
490    }
491
492    /// Convert a SPARQL triple pattern to a model triple pattern
493    fn convert_sparql_triple_pattern(
494        &self,
495        pattern: &SparqlTriplePattern,
496    ) -> Result<crate::model::pattern::TriplePattern, OxirsError> {
497        use crate::model::pattern::*;
498
499        let subject = match &pattern.subject {
500            SparqlTermPattern::Variable(v) => Some(SubjectPattern::Variable(v.clone())),
501            SparqlTermPattern::NamedNode(n) => Some(SubjectPattern::NamedNode(n.clone())),
502            SparqlTermPattern::BlankNode(b) => Some(SubjectPattern::BlankNode(b.clone())),
503            _ => None,
504        };
505
506        let predicate = match &pattern.predicate {
507            SparqlTermPattern::Variable(v) => Some(PredicatePattern::Variable(v.clone())),
508            SparqlTermPattern::NamedNode(n) => Some(PredicatePattern::NamedNode(n.clone())),
509            _ => None,
510        };
511
512        let object = match &pattern.object {
513            SparqlTermPattern::Variable(v) => Some(ObjectPattern::Variable(v.clone())),
514            SparqlTermPattern::NamedNode(n) => Some(ObjectPattern::NamedNode(n.clone())),
515            SparqlTermPattern::BlankNode(b) => Some(ObjectPattern::BlankNode(b.clone())),
516            SparqlTermPattern::Literal(l) => Some(ObjectPattern::Literal(l.clone())),
517            #[cfg(feature = "sparql-12")]
518            SparqlTermPattern::Triple(_) => {
519                // Triple patterns in object position not yet fully supported
520                None
521            }
522        };
523
524        Ok(crate::model::pattern::TriplePattern {
525            subject,
526            predicate,
527            object,
528        })
529    }
530
531    /// Convert a SPARQL algebra triple pattern to a model triple pattern
532    #[allow(dead_code)]
533    fn convert_triple_pattern(
534        &self,
535        pattern: &AlgebraTriplePattern,
536    ) -> Result<crate::model::pattern::TriplePattern, OxirsError> {
537        use crate::model::pattern::*;
538
539        let subject = match &pattern.subject {
540            TermPattern::Variable(v) => Some(SubjectPattern::Variable(v.clone())),
541            TermPattern::NamedNode(n) => Some(SubjectPattern::NamedNode(n.clone())),
542            TermPattern::BlankNode(b) => Some(SubjectPattern::BlankNode(b.clone())),
543            _ => None,
544        };
545
546        let predicate = match &pattern.predicate {
547            TermPattern::Variable(v) => Some(PredicatePattern::Variable(v.clone())),
548            TermPattern::NamedNode(n) => Some(PredicatePattern::NamedNode(n.clone())),
549            _ => None,
550        };
551
552        let object = match &pattern.object {
553            TermPattern::Variable(v) => Some(ObjectPattern::Variable(v.clone())),
554            TermPattern::NamedNode(n) => Some(ObjectPattern::NamedNode(n.clone())),
555            TermPattern::BlankNode(b) => Some(ObjectPattern::BlankNode(b.clone())),
556            TermPattern::Literal(l) => Some(ObjectPattern::Literal(l.clone())),
557            TermPattern::QuotedTriple(_) => {
558                return Err(OxirsError::Query(
559                    "RDF-star quoted triples are not yet supported as object patterns in the query module".to_string(),
560                ));
561            }
562        };
563
564        Ok(crate::model::pattern::TriplePattern {
565            subject,
566            predicate,
567            object,
568        })
569    }
570
571    /// Find variables that appear in (are produced by) both execution plans.
572    ///
573    /// These shared variables are the hash-join keys. Returning them lets the
574    /// executor bucket solutions by their common bindings instead of hashing
575    /// every solution under the empty key (which degenerates into a full
576    /// cartesian product materialised in memory).
577    fn find_join_variables(&self, left: &ExecutionPlan, right: &ExecutionPlan) -> Vec<Variable> {
578        let mut left_vars = std::collections::HashSet::new();
579        Self::collect_plan_variables(left, &mut left_vars);
580        let mut right_vars = std::collections::HashSet::new();
581        Self::collect_plan_variables(right, &mut right_vars);
582
583        // Deterministic ordering keeps the join key stable across runs.
584        let mut shared: Vec<Variable> = left_vars.intersection(&right_vars).cloned().collect();
585        shared.sort_by(|a, b| a.name().cmp(b.name()));
586        shared
587    }
588
589    /// Collect every variable an execution plan can bind in its output rows.
590    fn collect_plan_variables(
591        plan: &ExecutionPlan,
592        vars: &mut std::collections::HashSet<Variable>,
593    ) {
594        use crate::model::pattern::{ObjectPattern, PredicatePattern, SubjectPattern};
595        match plan {
596            ExecutionPlan::TripleScan { pattern } => {
597                if let Some(SubjectPattern::Variable(v)) = pattern.subject() {
598                    vars.insert(v.clone());
599                }
600                if let Some(PredicatePattern::Variable(v)) = pattern.predicate() {
601                    vars.insert(v.clone());
602                }
603                if let Some(ObjectPattern::Variable(v)) = pattern.object() {
604                    vars.insert(v.clone());
605                }
606            }
607            ExecutionPlan::HashJoin { left, right, .. } | ExecutionPlan::Union { left, right } => {
608                Self::collect_plan_variables(left, vars);
609                Self::collect_plan_variables(right, vars);
610            }
611            ExecutionPlan::Filter { input, .. }
612            | ExecutionPlan::Sort { input, .. }
613            | ExecutionPlan::Limit { input, .. }
614            | ExecutionPlan::Distinct { input } => {
615                Self::collect_plan_variables(input, vars);
616            }
617            ExecutionPlan::Project { vars: proj, .. } => {
618                vars.extend(proj.iter().cloned());
619            }
620        }
621    }
622
623    /// Convert a sparql_algebra ordering condition into an algebra
624    /// [`OrderExpression`] usable by the execution plan.
625    fn convert_order_expression(
626        &self,
627        order: &sparql_algebra::OrderExpression,
628    ) -> Result<crate::query::algebra::OrderExpression, OxirsError> {
629        use crate::query::algebra::OrderExpression as AlgebraOrder;
630        Ok(match order {
631            sparql_algebra::OrderExpression::Asc(expr) => {
632                AlgebraOrder::Asc(self.convert_expression(expr.clone())?)
633            }
634            sparql_algebra::OrderExpression::Desc(expr) => {
635                AlgebraOrder::Desc(self.convert_expression(expr.clone())?)
636            }
637        })
638    }
639
640    /// Convert sparql_algebra::Expression to algebra::Expression
641    #[allow(clippy::only_used_in_recursion)]
642    fn convert_expression(
643        &self,
644        expr: sparql_algebra::Expression,
645    ) -> Result<AlgebraExpression, OxirsError> {
646        use sparql_algebra::Expression as SparqlExpr;
647        use AlgebraExpression as AlgebraExpr;
648
649        match expr {
650            SparqlExpr::NamedNode(n) => Ok(AlgebraExpr::Term(crate::model::Term::NamedNode(n))),
651            SparqlExpr::Literal(l) => Ok(AlgebraExpr::Term(crate::model::Term::Literal(l))),
652            SparqlExpr::Variable(v) => Ok(AlgebraExpr::Variable(v)),
653            SparqlExpr::Or(left, right) => {
654                let left_expr = self.convert_expression(*left)?;
655                let right_expr = self.convert_expression(*right)?;
656                Ok(AlgebraExpr::Or(Box::new(left_expr), Box::new(right_expr)))
657            }
658            SparqlExpr::And(left, right) => {
659                let left_expr = self.convert_expression(*left)?;
660                let right_expr = self.convert_expression(*right)?;
661                Ok(AlgebraExpr::And(Box::new(left_expr), Box::new(right_expr)))
662            }
663            SparqlExpr::Equal(left, right) => {
664                let left_expr = self.convert_expression(*left)?;
665                let right_expr = self.convert_expression(*right)?;
666                Ok(AlgebraExpr::Equal(
667                    Box::new(left_expr),
668                    Box::new(right_expr),
669                ))
670            }
671            SparqlExpr::SameTerm(left, right) => {
672                let left_expr = self.convert_expression(*left)?;
673                let right_expr = self.convert_expression(*right)?;
674                Ok(AlgebraExpr::Equal(
675                    Box::new(left_expr),
676                    Box::new(right_expr),
677                )) // Map SameTerm to Equal for now
678            }
679            SparqlExpr::Greater(left, right) => {
680                let left_expr = self.convert_expression(*left)?;
681                let right_expr = self.convert_expression(*right)?;
682                Ok(AlgebraExpr::Greater(
683                    Box::new(left_expr),
684                    Box::new(right_expr),
685                ))
686            }
687            SparqlExpr::GreaterOrEqual(left, right) => {
688                let left_expr = self.convert_expression(*left)?;
689                let right_expr = self.convert_expression(*right)?;
690                Ok(AlgebraExpr::GreaterOrEqual(
691                    Box::new(left_expr),
692                    Box::new(right_expr),
693                ))
694            }
695            SparqlExpr::Less(left, right) => {
696                let left_expr = self.convert_expression(*left)?;
697                let right_expr = self.convert_expression(*right)?;
698                Ok(AlgebraExpr::Less(Box::new(left_expr), Box::new(right_expr)))
699            }
700            SparqlExpr::LessOrEqual(left, right) => {
701                let left_expr = self.convert_expression(*left)?;
702                let right_expr = self.convert_expression(*right)?;
703                Ok(AlgebraExpr::LessOrEqual(
704                    Box::new(left_expr),
705                    Box::new(right_expr),
706                ))
707            }
708            SparqlExpr::Not(inner) => {
709                let inner_expr = self.convert_expression(*inner)?;
710                Ok(AlgebraExpr::Not(Box::new(inner_expr)))
711            }
712            SparqlExpr::Bound(var) => Ok(AlgebraExpr::Bound(var)),
713            _ => {
714                // For expressions not yet supported, create a placeholder
715                Err(OxirsError::Query(format!(
716                    "Expression type not yet supported in conversion: {expr:?}"
717                )))
718            }
719        }
720    }
721
722    /// Determine the result variables of a SELECT query, honoring an explicit
723    /// projection when present (in SELECT order), otherwise falling back to all
724    /// variables bound by the pattern (sorted, `SELECT *` semantics).
725    fn result_variables(&self, pattern: &SparqlGraphPattern) -> Vec<Variable> {
726        if let Some(vars) = Self::projection_of(pattern) {
727            let mut seen = std::collections::HashSet::new();
728            vars.into_iter()
729                .filter(|v| seen.insert(v.clone()))
730                .collect()
731        } else {
732            self.extract_variables(pattern)
733        }
734    }
735
736    /// Find the projection variable list, descending through the solution
737    /// modifier wrappers (Distinct/Reduced/Slice/OrderBy) that sit above a
738    /// Project node.
739    fn projection_of(pattern: &SparqlGraphPattern) -> Option<Vec<Variable>> {
740        match pattern {
741            SparqlGraphPattern::Project { variables, .. } => Some(variables.clone()),
742            SparqlGraphPattern::Distinct { inner }
743            | SparqlGraphPattern::Reduced { inner }
744            | SparqlGraphPattern::Slice { inner, .. }
745            | SparqlGraphPattern::OrderBy { inner, .. } => Self::projection_of(inner),
746            _ => None,
747        }
748    }
749
750    /// Extract all variables from a graph pattern
751    fn extract_variables(&self, pattern: &SparqlGraphPattern) -> Vec<Variable> {
752        let mut variables = Vec::new();
753        self.collect_variables_from_pattern(pattern, &mut variables);
754        variables.sort_by_key(|v: &Variable| v.name().to_owned());
755        variables.dedup();
756        variables
757    }
758
759    /// Recursively collect variables from a graph pattern
760    fn collect_variables_from_pattern(
761        &self,
762        pattern: &SparqlGraphPattern,
763        variables: &mut Vec<Variable>,
764    ) {
765        match pattern {
766            SparqlGraphPattern::Bgp { patterns } => {
767                for triple_pattern in patterns {
768                    self.collect_variables_from_triple_pattern(triple_pattern, variables);
769                }
770            }
771            SparqlGraphPattern::Join { left, right } => {
772                self.collect_variables_from_pattern(left, variables);
773                self.collect_variables_from_pattern(right, variables);
774            }
775            SparqlGraphPattern::Filter { inner, .. } => {
776                self.collect_variables_from_pattern(inner, variables);
777            }
778            SparqlGraphPattern::Union { left, right } => {
779                self.collect_variables_from_pattern(left, variables);
780                self.collect_variables_from_pattern(right, variables);
781            }
782            SparqlGraphPattern::Project {
783                inner,
784                variables: proj_vars,
785            } => {
786                self.collect_variables_from_pattern(inner, variables);
787                variables.extend(proj_vars.iter().cloned());
788            }
789            SparqlGraphPattern::Distinct { inner } => {
790                self.collect_variables_from_pattern(inner, variables);
791            }
792            SparqlGraphPattern::Slice { inner, .. } => {
793                self.collect_variables_from_pattern(inner, variables);
794            }
795            _ => {
796                // Handle other pattern types as needed
797            }
798        }
799    }
800
801    /// Collect variables from a triple pattern
802    fn collect_variables_from_triple_pattern(
803        &self,
804        pattern: &SparqlTriplePattern,
805        variables: &mut Vec<Variable>,
806    ) {
807        if let SparqlTermPattern::Variable(v) = &pattern.subject {
808            variables.push(v.clone());
809        }
810        if let SparqlTermPattern::Variable(v) = &pattern.predicate {
811            variables.push(v.clone());
812        }
813        if let SparqlTermPattern::Variable(v) = &pattern.object {
814            variables.push(v.clone());
815        }
816    }
817
818    /// Instantiate a triple pattern with a solution
819    fn instantiate_triple_pattern(
820        &self,
821        pattern: &SparqlTriplePattern,
822        solution: &Solution,
823    ) -> Result<Option<crate::model::Triple>, OxirsError> {
824        use crate::model::*;
825
826        let subject = match &pattern.subject {
827            SparqlTermPattern::Variable(v) => {
828                if let Some(term) = solution.get(v) {
829                    match term {
830                        Term::NamedNode(n) => Subject::NamedNode(n.clone()),
831                        Term::BlankNode(b) => Subject::BlankNode(b.clone()),
832                        _ => return Ok(None), // Invalid subject
833                    }
834                } else {
835                    return Ok(None); // Unbound variable
836                }
837            }
838            SparqlTermPattern::NamedNode(n) => Subject::NamedNode(n.clone()),
839            SparqlTermPattern::BlankNode(b) => Subject::BlankNode(b.clone()),
840            _ => return Ok(None), // Invalid subject pattern
841        };
842
843        let predicate = match &pattern.predicate {
844            SparqlTermPattern::Variable(v) => {
845                if let Some(Term::NamedNode(n)) = solution.get(v) {
846                    Predicate::NamedNode(n.clone())
847                } else {
848                    return Ok(None); // Unbound or invalid predicate
849                }
850            }
851            SparqlTermPattern::NamedNode(n) => Predicate::NamedNode(n.clone()),
852            _ => return Ok(None), // Invalid predicate pattern
853        };
854
855        let object = match &pattern.object {
856            SparqlTermPattern::Variable(v) => {
857                if let Some(term) = solution.get(v) {
858                    match term {
859                        Term::NamedNode(n) => Object::NamedNode(n.clone()),
860                        Term::BlankNode(b) => Object::BlankNode(b.clone()),
861                        Term::Literal(l) => Object::Literal(l.clone()),
862                        _ => return Ok(None), // Invalid object
863                    }
864                } else {
865                    return Ok(None); // Unbound variable
866                }
867            }
868            SparqlTermPattern::NamedNode(n) => Object::NamedNode(n.clone()),
869            SparqlTermPattern::BlankNode(b) => Object::BlankNode(b.clone()),
870            SparqlTermPattern::Literal(l) => Object::Literal(l.clone()),
871            #[cfg(feature = "sparql-12")]
872            SparqlTermPattern::Triple(_) => {
873                // Triple patterns in object position not yet fully supported
874                return Ok(None);
875            }
876        };
877
878        Ok(Some(Triple::new(subject, predicate, object)))
879    }
880
881    /// Execute a SPARQL query string against a store (async version with federation support)
882    pub async fn query_async(
883        &self,
884        query_str: &str,
885        store: &dyn Store,
886    ) -> Result<QueryResult, OxirsError> {
887        // Parse the query string
888        let parsed_query = self.parser.parse(query_str)?;
889
890        // Execute the parsed query
891        self.execute_query_async(&parsed_query, store).await
892    }
893
894    /// Execute a parsed Query object against a store (async version)
895    pub async fn execute_query_async(
896        &self,
897        query: &sparql_query::Query,
898        store: &dyn Store,
899    ) -> Result<QueryResult, OxirsError> {
900        match query {
901            sparql_query::Query::Select {
902                pattern, dataset, ..
903            } => {
904                self.execute_select_query_async(pattern, dataset.as_ref(), store)
905                    .await
906            }
907            sparql_query::Query::Ask {
908                pattern, dataset, ..
909            } => {
910                self.execute_ask_query_async(pattern, dataset.as_ref(), store)
911                    .await
912            }
913            sparql_query::Query::Construct {
914                template,
915                pattern,
916                dataset,
917                ..
918            } => {
919                self.execute_construct_query_async(template, pattern, dataset.as_ref(), store)
920                    .await
921            }
922            sparql_query::Query::Describe {
923                pattern, dataset, ..
924            } => {
925                self.execute_describe_query_async(pattern, dataset.as_ref(), store)
926                    .await
927            }
928        }
929    }
930
931    /// Execute a SELECT query (async version with federation support)
932    async fn execute_select_query_async(
933        &self,
934        pattern: &SparqlGraphPattern,
935        _dataset: Option<&QueryDataset>,
936        store: &dyn Store,
937    ) -> Result<QueryResult, OxirsError> {
938        // Check if pattern contains SERVICE clause
939        if self.contains_service_clause(pattern) {
940            // Use federation executor
941            return self.execute_federated_select(pattern, store).await;
942        }
943
944        // Fall back to regular execution
945        self.execute_select_query(pattern, _dataset, store)
946    }
947
948    /// Execute an ASK query (async version)
949    async fn execute_ask_query_async(
950        &self,
951        pattern: &SparqlGraphPattern,
952        dataset: Option<&QueryDataset>,
953        store: &dyn Store,
954    ) -> Result<QueryResult, OxirsError> {
955        if self.contains_service_clause(pattern) {
956            let result = self.execute_federated_select(pattern, store).await?;
957            if let QueryResult::Select { bindings, .. } = result {
958                return Ok(QueryResult::Ask(!bindings.is_empty()));
959            }
960        }
961        self.execute_ask_query(pattern, dataset, store)
962    }
963
964    /// Execute a CONSTRUCT query (async version)
965    async fn execute_construct_query_async(
966        &self,
967        template: &[SparqlTriplePattern],
968        pattern: &SparqlGraphPattern,
969        dataset: Option<&QueryDataset>,
970        store: &dyn Store,
971    ) -> Result<QueryResult, OxirsError> {
972        if self.contains_service_clause(pattern) {
973            return Err(OxirsError::Federation(
974                "CONSTRUCT with SERVICE is not yet fully supported".to_string(),
975            ));
976        }
977        self.execute_construct_query(template, pattern, dataset, store)
978    }
979
980    /// Execute a DESCRIBE query (async version)
981    async fn execute_describe_query_async(
982        &self,
983        pattern: &SparqlGraphPattern,
984        dataset: Option<&QueryDataset>,
985        store: &dyn Store,
986    ) -> Result<QueryResult, OxirsError> {
987        if self.contains_service_clause(pattern) {
988            return Err(OxirsError::Federation(
989                "DESCRIBE with SERVICE is not yet fully supported".to_string(),
990            ));
991        }
992        self.execute_describe_query(pattern, dataset, store)
993    }
994
995    /// Check if a pattern contains a SERVICE clause
996    fn contains_service_clause(&self, pattern: &SparqlGraphPattern) -> bool {
997        // Check for SERVICE patterns using recursive pattern traversal
998        matches!(pattern, SparqlGraphPattern::Service { .. })
999            || Self::pattern_contains_service_recursive(pattern)
1000    }
1001
1002    /// Recursively check for SERVICE in nested patterns
1003    fn pattern_contains_service_recursive(pattern: &SparqlGraphPattern) -> bool {
1004        match pattern {
1005            SparqlGraphPattern::Service { .. } => true,
1006            SparqlGraphPattern::Join { left, right }
1007            | SparqlGraphPattern::Union { left, right } => {
1008                Self::pattern_contains_service_recursive(left)
1009                    || Self::pattern_contains_service_recursive(right)
1010            }
1011            SparqlGraphPattern::Filter { inner, .. }
1012            | SparqlGraphPattern::Distinct { inner }
1013            | SparqlGraphPattern::Reduced { inner }
1014            | SparqlGraphPattern::Project { inner, .. } => {
1015                Self::pattern_contains_service_recursive(inner)
1016            }
1017            SparqlGraphPattern::LeftJoin { left, right, .. } => {
1018                Self::pattern_contains_service_recursive(left)
1019                    || Self::pattern_contains_service_recursive(right)
1020            }
1021            _ => false,
1022        }
1023    }
1024
1025    /// Execute a federated SELECT query
1026    async fn execute_federated_select(
1027        &self,
1028        pattern: &SparqlGraphPattern,
1029        store: &dyn Store,
1030    ) -> Result<QueryResult, OxirsError> {
1031        let federation_executor = self.federation_executor.as_ref().ok_or_else(|| {
1032            OxirsError::Federation("Federation executor not available".to_string())
1033        })?;
1034
1035        // Start with empty bindings from local store
1036        let local_bindings = Vec::new();
1037
1038        // Execute the federated pattern
1039        let bindings = self
1040            .execute_pattern_with_federation(pattern, local_bindings, federation_executor, store)
1041            .await?;
1042
1043        // Extract result variable names (honoring any explicit projection)
1044        let variables = self
1045            .result_variables(pattern)
1046            .into_iter()
1047            .map(|v| v.name().to_string())
1048            .collect();
1049
1050        Ok(QueryResult::Select {
1051            variables,
1052            bindings,
1053        })
1054    }
1055
1056    /// Execute a pattern with federation support
1057    fn execute_pattern_with_federation<'a>(
1058        &'a self,
1059        pattern: &'a SparqlGraphPattern,
1060        current_bindings: Vec<HashMap<String, Term>>,
1061        federation_executor: &'a crate::federation::FederationExecutor,
1062        store: &'a dyn Store,
1063    ) -> FederatedQueryFuture<'a> {
1064        Box::pin(async move {
1065            let current_bindings = current_bindings;
1066            match pattern {
1067                SparqlGraphPattern::Service {
1068                    name,
1069                    inner,
1070                    silent,
1071                } => {
1072                    // Execute SERVICE clause
1073                    let remote_bindings = federation_executor
1074                        .execute_service(name, inner, *silent, &current_bindings)
1075                        .await?;
1076
1077                    // Merge local and remote bindings
1078                    if current_bindings.is_empty() {
1079                        Ok(remote_bindings)
1080                    } else {
1081                        Ok(federation_executor.merge_bindings(current_bindings, remote_bindings))
1082                    }
1083                }
1084                SparqlGraphPattern::Join { left, right } => {
1085                    // Execute left side first
1086                    let left_bindings = self
1087                        .execute_pattern_with_federation(
1088                            left,
1089                            current_bindings,
1090                            federation_executor,
1091                            store,
1092                        )
1093                        .await?;
1094
1095                    // Then execute right side with left bindings
1096                    self.execute_pattern_with_federation(
1097                        right,
1098                        left_bindings,
1099                        federation_executor,
1100                        store,
1101                    )
1102                    .await
1103                }
1104                _ => {
1105                    // For non-federated patterns, use regular execution
1106                    let executor = QueryExecutor::new(store);
1107                    let plan = self.pattern_to_plan(pattern)?;
1108                    let solutions = executor.execute(&plan)?;
1109
1110                    let bindings: Vec<HashMap<String, Term>> = solutions
1111                        .into_iter()
1112                        .take(self.executor_config.max_results)
1113                        .map(|sol| {
1114                            sol.iter()
1115                                .map(|(var, term)| (var.name().to_string(), term.clone()))
1116                                .collect()
1117                        })
1118                        .collect();
1119
1120                    if current_bindings.is_empty() {
1121                        Ok(bindings)
1122                    } else {
1123                        // Merge current and new bindings
1124                        Ok(federation_executor.merge_bindings(current_bindings, bindings))
1125                    }
1126                }
1127            }
1128        })
1129    }
1130
1131    /// Check if a triple involves a specific term
1132    fn triple_involves_term(&self, triple: &crate::model::Triple, term: &Term) -> bool {
1133        match term {
1134            Term::NamedNode(n) => {
1135                matches!(triple.subject(), Subject::NamedNode(sn) if sn == n)
1136                    || matches!(triple.predicate(), Predicate::NamedNode(pn) if pn == n)
1137                    || matches!(triple.object(), Object::NamedNode(on) if on == n)
1138            }
1139            Term::BlankNode(b) => {
1140                matches!(triple.subject(), Subject::BlankNode(sb) if sb == b)
1141                    || matches!(triple.object(), Object::BlankNode(ob) if ob == b)
1142            }
1143            Term::Literal(l) => {
1144                matches!(triple.object(), Object::Literal(ol) if ol == l)
1145            }
1146            _ => false,
1147        }
1148    }
1149}
1150
1151#[cfg(test)]
1152mod join_variable_tests {
1153    use super::*;
1154    use crate::model::pattern::{ObjectPattern, PredicatePattern, SubjectPattern, TriplePattern};
1155
1156    fn scan(subject: &str, predicate: &str, object: &str) -> ExecutionPlan {
1157        ExecutionPlan::TripleScan {
1158            pattern: TriplePattern::new(
1159                Some(SubjectPattern::Variable(
1160                    Variable::new(subject).expect("var"),
1161                )),
1162                Some(PredicatePattern::NamedNode(
1163                    crate::model::NamedNode::new(predicate).expect("iri"),
1164                )),
1165                Some(ObjectPattern::Variable(Variable::new(object).expect("var"))),
1166            ),
1167        }
1168    }
1169
1170    /// P1: find_join_variables must return the variables shared by both plans so
1171    /// the hash join keys correctly instead of degenerating into a cartesian
1172    /// product hashed under the empty key.
1173    #[test]
1174    fn regression_find_join_variables_returns_shared() {
1175        let engine = QueryEngine::new();
1176        // { ?x :p ?y } JOIN { ?y :q ?z } shares ?y.
1177        let left = scan("?x", "http://example.org/p", "?y");
1178        let right = scan("?y", "http://example.org/q", "?z");
1179
1180        let join_vars = engine.find_join_variables(&left, &right);
1181        assert_eq!(join_vars.len(), 1, "exactly one shared variable expected");
1182        assert_eq!(join_vars[0].name(), "y");
1183    }
1184
1185    /// When the two plans share no variable the join key must be empty (a real
1186    /// cartesian product, correctly represented).
1187    #[test]
1188    fn regression_find_join_variables_disjoint_is_empty() {
1189        let engine = QueryEngine::new();
1190        let left = scan("?a", "http://example.org/p", "?b");
1191        let right = scan("?c", "http://example.org/q", "?d");
1192        assert!(engine.find_join_variables(&left, &right).is_empty());
1193    }
1194}