Skip to main content

oxirs_arq/execution/
parallel_eval.rs

1//! Parallel Triple Pattern Evaluation
2//!
3//! This module implements parallel evaluation of independent triple patterns
4//! within a Basic Graph Pattern (BGP).  Independent patterns (those that do
5//! not share variables with each other) can be executed concurrently, and
6//! their results are then joined in the correct order.
7//!
8//! The dependency analysis uses a directed acyclic graph (DAG) over patterns
9//! to identify which patterns can be safely parallelised.
10
11use crate::optimizer::adaptive::TriplePatternInfo;
12use crate::optimizer::materialized_view::{BindingRow, RdfTerm};
13use anyhow::Result;
14use std::collections::{HashMap, HashSet, VecDeque};
15use std::sync::Arc;
16
17// ---------------------------------------------------------------------------
18// Public trait: TripleStore
19// ---------------------------------------------------------------------------
20
21/// Trait abstracting a triple store that can evaluate a single pattern
22/// against an optional set of input bindings.
23///
24/// Implementations must be `Send + Sync` to allow parallel evaluation.
25pub trait TripleStore: Send + Sync {
26    /// Evaluate a triple pattern, optionally constrained by `bindings`.
27    ///
28    /// If `bindings` is `Some`, each row is used to ground variables before
29    /// evaluating the pattern (index-nested-loop style).
30    fn evaluate_pattern(
31        &self,
32        pattern: &TriplePatternInfo,
33        bindings: Option<&[BindingRow]>,
34    ) -> Result<Vec<BindingRow>>;
35
36    /// Return a quick cardinality estimate for a pattern (used for ordering).
37    fn estimate_cardinality(&self, pattern: &TriplePatternInfo) -> u64;
38}
39
40// ---------------------------------------------------------------------------
41// Dependency analysis
42// ---------------------------------------------------------------------------
43
44/// Dependency graph over a set of triple patterns.
45///
46/// Pattern `i` depends on pattern `j` when `j` binds a variable that `i`
47/// needs as input.  In the context of parallel evaluation, two patterns are
48/// *independent* when neither depends on the other.
49pub struct PatternDependencyGraph {
50    patterns: Vec<TriplePatternInfo>,
51    /// `dependencies[i]` = set of pattern indices that `i` depends on
52    dependencies: Vec<HashSet<usize>>,
53    /// Topologically sorted execution stages: each stage is a set of indices
54    /// that can be evaluated in parallel once all previous stages complete.
55    execution_stages: Vec<Vec<usize>>,
56}
57
58impl PatternDependencyGraph {
59    /// Build the dependency graph for the given pattern list.
60    pub fn build(patterns: Vec<TriplePatternInfo>) -> Self {
61        let n = patterns.len();
62        let mut dependencies: Vec<HashSet<usize>> = vec![HashSet::new(); n];
63
64        // Build a map from variable name -> first pattern that binds it
65        let mut var_producer: HashMap<String, usize> = HashMap::new();
66        for (i, pattern) in patterns.iter().enumerate() {
67            for var in &pattern.bound_variables {
68                var_producer.entry(var.clone()).or_insert(i);
69            }
70        }
71
72        // Pattern `i` depends on `j` if `j` is the producer of a variable
73        // that `i` also uses, and j != i.
74        for i in 0..n {
75            for (var_name, &producer) in &var_producer {
76                if producer == i {
77                    continue;
78                }
79                if patterns[i].bound_variables.contains(var_name) {
80                    dependencies[i].insert(producer);
81                }
82            }
83        }
84
85        let execution_stages = Self::topological_stages(&dependencies, n);
86
87        Self {
88            patterns,
89            dependencies,
90            execution_stages,
91        }
92    }
93
94    /// Return groups of patterns that can be evaluated in parallel.
95    pub fn get_independent_patterns(&self) -> Vec<Vec<usize>> {
96        self.execution_stages.clone()
97    }
98
99    /// Return the topologically sorted execution stages.
100    pub fn execution_order(&self) -> &[Vec<usize>] {
101        &self.execution_stages
102    }
103
104    /// Access the underlying patterns
105    pub fn patterns(&self) -> &[TriplePatternInfo] {
106        &self.patterns
107    }
108
109    /// Check whether two pattern indices are independent
110    pub fn are_independent(&self, i: usize, j: usize) -> bool {
111        !self.dependencies[i].contains(&j) && !self.dependencies[j].contains(&i)
112    }
113
114    // -----------------------------------------------------------------------
115    // Private helpers
116    // -----------------------------------------------------------------------
117
118    /// Kahn's algorithm for topological layering
119    fn topological_stages(dependencies: &[HashSet<usize>], n: usize) -> Vec<Vec<usize>> {
120        let mut in_degree: Vec<usize> = dependencies.iter().map(|d| d.len()).collect();
121        let mut reverse: Vec<Vec<usize>> = vec![Vec::new(); n];
122        for (i, deps) in dependencies.iter().enumerate() {
123            for &dep in deps {
124                reverse[dep].push(i);
125            }
126        }
127
128        let mut stages: Vec<Vec<usize>> = Vec::new();
129        let mut queue: VecDeque<usize> = in_degree
130            .iter()
131            .enumerate()
132            .filter(|(_, &d)| d == 0)
133            .map(|(i, _)| i)
134            .collect();
135
136        while !queue.is_empty() {
137            let stage: Vec<usize> = queue.drain(..).collect();
138            for &node in &stage {
139                for &dependent in &reverse[node] {
140                    in_degree[dependent] -= 1;
141                    if in_degree[dependent] == 0 {
142                        queue.push_back(dependent);
143                    }
144                }
145            }
146            stages.push(stage);
147        }
148
149        stages
150    }
151}
152
153// ---------------------------------------------------------------------------
154// Parallel BGP evaluator
155// ---------------------------------------------------------------------------
156
157/// Parallel evaluator for Basic Graph Patterns.
158pub struct ParallelBgpEvaluator {
159    /// Number of worker threads (0 = use Rayon's global pool)
160    pub num_threads: usize,
161    /// Minimum patterns per stage to justify parallelism overhead
162    pub chunk_size: usize,
163}
164
165impl Default for ParallelBgpEvaluator {
166    fn default() -> Self {
167        Self {
168            num_threads: std::thread::available_parallelism()
169                .map(|n| n.get())
170                .unwrap_or(1),
171            chunk_size: 1,
172        }
173    }
174}
175
176impl ParallelBgpEvaluator {
177    /// Create a new evaluator with a specific thread count
178    pub fn new(num_threads: usize) -> Self {
179        Self {
180            num_threads,
181            chunk_size: 1,
182        }
183    }
184
185    /// Create a new evaluator with tunable chunk size
186    pub fn with_chunk_size(mut self, chunk_size: usize) -> Self {
187        self.chunk_size = chunk_size.max(1);
188        self
189    }
190
191    /// Evaluate a BGP by exploiting parallelism among independent patterns.
192    pub fn evaluate(
193        &self,
194        patterns: Vec<TriplePatternInfo>,
195        store: &dyn TripleStore,
196    ) -> Result<Vec<BindingRow>> {
197        if patterns.is_empty() {
198            return Ok(Vec::new());
199        }
200
201        let graph = PatternDependencyGraph::build(patterns);
202        let stages = graph.execution_order().to_vec();
203
204        // Running binding set: starts as a single empty row (identity for join)
205        let mut current_bindings: Vec<BindingRow> = vec![BindingRow::new()];
206
207        for stage in &stages {
208            let stage_results =
209                self.evaluate_stage(stage, graph.patterns(), store, &current_bindings)?;
210
211            for (pattern_idx, pattern_rows) in stage_results {
212                let pattern = &graph.patterns()[pattern_idx];
213                // Only join on variables that already appear in current_bindings
214                let join_vars: Vec<String> = if current_bindings.is_empty() {
215                    Vec::new()
216                } else {
217                    let first_row = &current_bindings[0];
218                    pattern
219                        .bound_variables
220                        .iter()
221                        .filter(|v| first_row.contains_key(v.as_str()))
222                        .cloned()
223                        .collect()
224                };
225
226                current_bindings = self.merge_results(current_bindings, pattern_rows, &join_vars);
227            }
228        }
229
230        Ok(current_bindings)
231    }
232
233    /// Evaluate a set of pattern indices in parallel within a single stage.
234    fn evaluate_stage(
235        &self,
236        stage: &[usize],
237        patterns: &[TriplePatternInfo],
238        store: &dyn TripleStore,
239        current_bindings: &[BindingRow],
240    ) -> Result<Vec<(usize, Vec<BindingRow>)>> {
241        if stage.is_empty() {
242            return Ok(Vec::new());
243        }
244
245        if stage.len() < self.chunk_size || self.num_threads <= 1 {
246            return self.evaluate_stage_sequential(stage, patterns, store, current_bindings);
247        }
248
249        #[cfg(feature = "parallel")]
250        {
251            self.evaluate_stage_parallel(stage, patterns, store, current_bindings)
252        }
253        #[cfg(not(feature = "parallel"))]
254        {
255            self.evaluate_stage_sequential(stage, patterns, store, current_bindings)
256        }
257    }
258
259    fn evaluate_stage_sequential(
260        &self,
261        stage: &[usize],
262        patterns: &[TriplePatternInfo],
263        store: &dyn TripleStore,
264        current_bindings: &[BindingRow],
265    ) -> Result<Vec<(usize, Vec<BindingRow>)>> {
266        let mut results = Vec::with_capacity(stage.len());
267        for &idx in stage {
268            let rows = store.evaluate_pattern(&patterns[idx], Some(current_bindings))?;
269            results.push((idx, rows));
270        }
271        Ok(results)
272    }
273
274    #[cfg(feature = "parallel")]
275    fn evaluate_stage_parallel(
276        &self,
277        stage: &[usize],
278        patterns: &[TriplePatternInfo],
279        store: &dyn TripleStore,
280        current_bindings: &[BindingRow],
281    ) -> Result<Vec<(usize, Vec<BindingRow>)>> {
282        use rayon::prelude::*;
283        use std::sync::Mutex;
284
285        let error_cell: Arc<Mutex<Option<anyhow::Error>>> = Arc::new(Mutex::new(None));
286        let error_clone = Arc::clone(&error_cell);
287
288        let results: Vec<(usize, Vec<BindingRow>)> = stage
289            .par_iter()
290            .filter_map(|&idx| {
291                match store.evaluate_pattern(&patterns[idx], Some(current_bindings)) {
292                    Ok(rows) => Some((idx, rows)),
293                    Err(e) => {
294                        if let Ok(mut guard) = error_clone.lock() {
295                            if guard.is_none() {
296                                *guard = Some(e);
297                            }
298                        }
299                        None
300                    }
301                }
302            })
303            .collect();
304
305        if let Ok(mut guard) = error_cell.lock() {
306            if let Some(err) = guard.take() {
307                return Err(err);
308            }
309        }
310        Ok(results)
311    }
312
313    /// Hash join of two binding sets on shared variables.
314    pub fn merge_results(
315        &self,
316        left: Vec<BindingRow>,
317        right: Vec<BindingRow>,
318        join_vars: &[String],
319    ) -> Vec<BindingRow> {
320        if right.is_empty() {
321            return left;
322        }
323        if left.is_empty() {
324            return right;
325        }
326
327        // Cross product when no shared variables
328        if join_vars.is_empty() {
329            let mut output: Vec<BindingRow> = Vec::with_capacity(left.len() * right.len());
330            for l_row in &left {
331                for r_row in &right {
332                    let mut merged: BindingRow = l_row.clone();
333                    for (k, v) in r_row {
334                        merged.insert(k.clone(), v.clone());
335                    }
336                    output.push(merged);
337                }
338            }
339            return output;
340        }
341
342        // Build hash index over right side
343        let mut hash_index: HashMap<Vec<String>, Vec<usize>> = HashMap::new();
344        for (ridx, row) in right.iter().enumerate() {
345            let key: Vec<String> = join_vars.iter().map(|v| rdf_term_key(row.get(v))).collect();
346            hash_index.entry(key).or_default().push(ridx);
347        }
348
349        // Probe phase
350        let mut output: Vec<BindingRow> = Vec::new();
351        for l_row in &left {
352            let key: Vec<String> = join_vars
353                .iter()
354                .map(|v| rdf_term_key(l_row.get(v)))
355                .collect();
356
357            if let Some(right_indices) = hash_index.get(&key) {
358                for &ridx in right_indices {
359                    let r_row = &right[ridx];
360                    let mut merged: BindingRow = l_row.clone();
361                    for (k, v) in r_row {
362                        merged.insert(k.clone(), v.clone());
363                    }
364                    output.push(merged);
365                }
366            }
367        }
368        output
369    }
370}
371
372/// Convert an optional RdfTerm to a stable string key for hashing
373fn rdf_term_key(term: Option<&RdfTerm>) -> String {
374    match term {
375        None => String::new(),
376        Some(t) => format!("{t}"),
377    }
378}
379
380// ---------------------------------------------------------------------------
381// Mock TripleStore for testing
382// ---------------------------------------------------------------------------
383
384#[cfg(test)]
385pub(crate) mod test_support {
386    use super::*;
387
388    /// Simple in-memory triple store for tests.
389    pub struct MockTripleStore {
390        pub results: HashMap<String, Vec<BindingRow>>,
391        pub default_result: Vec<BindingRow>,
392    }
393
394    impl MockTripleStore {
395        pub fn new() -> Self {
396            Self {
397                results: HashMap::new(),
398                default_result: Vec::new(),
399            }
400        }
401
402        pub fn with_result(mut self, pattern_id: &str, rows: Vec<BindingRow>) -> Self {
403            self.results.insert(pattern_id.to_string(), rows);
404            self
405        }
406    }
407
408    impl TripleStore for MockTripleStore {
409        fn evaluate_pattern(
410            &self,
411            pattern: &TriplePatternInfo,
412            _bindings: Option<&[BindingRow]>,
413        ) -> Result<Vec<BindingRow>> {
414            Ok(self
415                .results
416                .get(&pattern.id)
417                .cloned()
418                .unwrap_or_else(|| self.default_result.clone()))
419        }
420
421        fn estimate_cardinality(&self, pattern: &TriplePatternInfo) -> u64 {
422            self.results
423                .get(&pattern.id)
424                .map(|r| r.len() as u64)
425                .unwrap_or(0)
426        }
427    }
428
429    pub fn iri_term(value: &str) -> RdfTerm {
430        RdfTerm::Iri(value.to_string())
431    }
432
433    pub fn make_row(pairs: &[(&str, RdfTerm)]) -> BindingRow {
434        pairs
435            .iter()
436            .map(|(k, v)| (k.to_string(), v.clone()))
437            .collect()
438    }
439}
440
441// ---------------------------------------------------------------------------
442// Tests
443// ---------------------------------------------------------------------------
444
445#[cfg(test)]
446mod tests {
447    use super::test_support::*;
448    use super::*;
449    use crate::optimizer::adaptive::{PatternTerm, TriplePatternInfo};
450    use crate::optimizer::materialized_view::RdfTerm;
451
452    fn simple_pattern(id: &str, vars: Vec<String>, cardinality: u64) -> TriplePatternInfo {
453        TriplePatternInfo {
454            id: id.to_string(),
455            subject: PatternTerm::Variable(vars.first().cloned().unwrap_or_default()),
456            predicate: PatternTerm::Iri(format!("http://example.org/p_{id}")),
457            object: PatternTerm::Variable(vars.last().cloned().unwrap_or_default()),
458            estimated_cardinality: cardinality,
459            bound_variables: vars,
460            original_pattern: None,
461        }
462    }
463
464    #[test]
465    fn test_dependency_graph_independent_patterns() {
466        let p1 = simple_pattern("p1", vec!["a".to_string(), "b".to_string()], 10);
467        let p2 = simple_pattern("p2", vec!["c".to_string(), "d".to_string()], 20);
468        let graph = PatternDependencyGraph::build(vec![p1, p2]);
469
470        assert!(
471            graph.are_independent(0, 1),
472            "Patterns with no shared vars should be independent"
473        );
474        let stages = graph.get_independent_patterns();
475        assert_eq!(stages.len(), 1, "Independent patterns fit into one stage");
476        assert_eq!(stages[0].len(), 2);
477    }
478
479    #[test]
480    fn test_dependency_graph_dependent_patterns() {
481        let p1 = simple_pattern("p1", vec!["s".to_string(), "type".to_string()], 10);
482        let p2 = simple_pattern("p2", vec!["s".to_string(), "name".to_string()], 100);
483        let graph = PatternDependencyGraph::build(vec![p1, p2]);
484
485        let stages = graph.get_independent_patterns();
486        let total: usize = stages.iter().map(|s| s.len()).sum();
487        assert_eq!(total, 2, "All patterns should appear across stages");
488    }
489
490    #[test]
491    fn test_parallel_evaluator_empty_patterns() {
492        let evaluator = ParallelBgpEvaluator::new(2);
493        let store = MockTripleStore::new();
494        let result = evaluator.evaluate(vec![], &store).unwrap();
495        assert!(result.is_empty());
496    }
497
498    #[test]
499    fn test_parallel_evaluator_single_pattern() {
500        let pattern = simple_pattern("pat1", vec!["s".to_string()], 2);
501        let rows = vec![
502            make_row(&[("s", iri_term("http://example.org/a"))]),
503            make_row(&[("s", iri_term("http://example.org/b"))]),
504        ];
505        let store = MockTripleStore::new().with_result("pat1", rows);
506        let evaluator = ParallelBgpEvaluator::new(1);
507        let result = evaluator.evaluate(vec![pattern], &store).unwrap();
508        assert_eq!(result.len(), 2);
509    }
510
511    #[test]
512    fn test_parallel_evaluator_two_patterns_with_join() {
513        let p1 = simple_pattern("p1", vec!["s".to_string(), "type".to_string()], 2);
514        let p2 = simple_pattern("p2", vec!["s".to_string(), "name".to_string()], 2);
515
516        let p1_rows = vec![
517            make_row(&[
518                ("s", iri_term("http://example.org/alice")),
519                ("type", iri_term("http://example.org/Person")),
520            ]),
521            make_row(&[
522                ("s", iri_term("http://example.org/bob")),
523                ("type", iri_term("http://example.org/Person")),
524            ]),
525        ];
526        let p2_rows = vec![
527            make_row(&[
528                ("s", iri_term("http://example.org/alice")),
529                ("name", RdfTerm::plain_literal("Alice")),
530            ]),
531            make_row(&[
532                ("s", iri_term("http://example.org/bob")),
533                ("name", RdfTerm::plain_literal("Bob")),
534            ]),
535        ];
536
537        let store = MockTripleStore::new()
538            .with_result("p1", p1_rows)
539            .with_result("p2", p2_rows);
540
541        let evaluator = ParallelBgpEvaluator::new(2);
542        let result = evaluator.evaluate(vec![p1, p2], &store).unwrap();
543
544        assert_eq!(
545            result.len(),
546            2,
547            "Should produce 2 joined rows (one per person)"
548        );
549        for row in &result {
550            assert!(row.contains_key("s"));
551            assert!(row.contains_key("name"));
552        }
553    }
554
555    #[test]
556    fn test_merge_results_no_join_vars_cross_product() {
557        let evaluator = ParallelBgpEvaluator::new(1);
558        let left = vec![
559            make_row(&[("a", iri_term("http://example.org/1"))]),
560            make_row(&[("a", iri_term("http://example.org/2"))]),
561        ];
562        let right = vec![make_row(&[("b", iri_term("http://example.org/x"))])];
563
564        let merged = evaluator.merge_results(left, right, &[]);
565        assert_eq!(merged.len(), 2, "Cross product of 2x1 = 2 rows");
566    }
567
568    #[test]
569    fn test_merge_results_with_join_var() {
570        let evaluator = ParallelBgpEvaluator::new(1);
571        let left = vec![
572            make_row(&[
573                ("s", iri_term("http://a")),
574                ("type", iri_term("http://Person")),
575            ]),
576            make_row(&[
577                ("s", iri_term("http://b")),
578                ("type", iri_term("http://Person")),
579            ]),
580        ];
581        let right = vec![make_row(&[
582            ("s", iri_term("http://a")),
583            ("name", RdfTerm::plain_literal("Alice")),
584        ])];
585
586        let merged = evaluator.merge_results(left, right, &["s".to_string()]);
587        assert_eq!(merged.len(), 1);
588        assert_eq!(
589            merged[0].get("name"),
590            Some(&RdfTerm::plain_literal("Alice"))
591        );
592    }
593
594    #[test]
595    fn test_merge_results_empty_left_returns_right() {
596        let evaluator = ParallelBgpEvaluator::new(1);
597        let right = vec![make_row(&[("s", iri_term("http://a"))])];
598        let merged = evaluator.merge_results(vec![], right, &[]);
599        assert_eq!(merged.len(), 1);
600    }
601
602    #[test]
603    fn test_merge_results_empty_right_returns_left() {
604        let evaluator = ParallelBgpEvaluator::new(1);
605        let left = vec![make_row(&[("s", iri_term("http://a"))])];
606        let merged = evaluator.merge_results(left, vec![], &[]);
607        assert_eq!(merged.len(), 1);
608    }
609
610    #[test]
611    fn test_dependency_graph_three_chain() {
612        let p1 = simple_pattern("p1", vec!["x".to_string()], 5);
613        let p2 = simple_pattern("p2", vec!["x".to_string(), "y".to_string()], 50);
614        let p3 = simple_pattern("p3", vec!["y".to_string(), "z".to_string()], 500);
615
616        let graph = PatternDependencyGraph::build(vec![p1, p2, p3]);
617        let stages = graph.execution_order();
618        let total: usize = stages.iter().map(|s| s.len()).sum();
619        assert_eq!(total, 3);
620        assert!(!stages.is_empty());
621    }
622
623    #[test]
624    fn test_evaluator_default_thread_count() {
625        let evaluator = ParallelBgpEvaluator::default();
626        assert!(evaluator.num_threads >= 1);
627    }
628}
629
630#[cfg(test)]
631mod extended_tests {
632    use super::test_support::*;
633    use super::*;
634    use crate::optimizer::adaptive::{PatternTerm, TriplePatternInfo};
635    use crate::optimizer::materialized_view::RdfTerm;
636
637    fn pat(id: &str, vars: Vec<String>, cardinality: u64) -> TriplePatternInfo {
638        TriplePatternInfo {
639            id: id.to_string(),
640            subject: PatternTerm::Variable(vars.first().cloned().unwrap_or_default()),
641            predicate: PatternTerm::Iri(format!("http://example.org/p_{id}")),
642            object: PatternTerm::Variable(vars.last().cloned().unwrap_or_default()),
643            estimated_cardinality: cardinality,
644            bound_variables: vars,
645            original_pattern: None,
646        }
647    }
648
649    // --- PatternDependencyGraph extended tests ---
650
651    #[test]
652    fn test_dependency_graph_single_pattern() {
653        let p1 = pat("solo", vec!["x".to_string()], 10);
654        let graph = PatternDependencyGraph::build(vec![p1]);
655
656        let stages = graph.get_independent_patterns();
657        assert_eq!(
658            stages.len(),
659            1,
660            "Single pattern should produce a single stage"
661        );
662        assert_eq!(stages[0], vec![0], "Stage 0 should contain pattern 0");
663    }
664
665    #[test]
666    fn test_dependency_graph_no_patterns() {
667        let graph = PatternDependencyGraph::build(vec![]);
668        assert!(graph.get_independent_patterns().is_empty());
669    }
670
671    #[test]
672    fn test_dependency_graph_are_independent_different_vars() {
673        let p1 = pat("p1", vec!["a".to_string(), "b".to_string()], 10);
674        let p2 = pat("p2", vec!["c".to_string(), "d".to_string()], 10);
675        let graph = PatternDependencyGraph::build(vec![p1, p2]);
676
677        assert!(
678            graph.are_independent(0, 1),
679            "Patterns with disjoint variables should be independent"
680        );
681    }
682
683    #[test]
684    fn test_dependency_graph_are_not_independent_shared_var() {
685        let p1 = pat("p1", vec!["s".to_string(), "o1".to_string()], 10);
686        let p2 = pat("p2", vec!["s".to_string(), "o2".to_string()], 10);
687        let graph = PatternDependencyGraph::build(vec![p1, p2]);
688
689        // The graph.are_independent() checks if both directions are dependency-free
690        // Shared variable "s" should create a dependency
691        // They can still be in the same stage if neither depends on the other's bindings
692        let _stages = graph.get_independent_patterns();
693        let patterns = graph.patterns();
694        assert_eq!(patterns.len(), 2, "Graph should contain 2 patterns");
695    }
696
697    #[test]
698    fn test_dependency_graph_execution_order_returns_all_patterns() {
699        let p1 = pat("p1", vec!["a".to_string()], 10);
700        let p2 = pat("p2", vec!["b".to_string()], 20);
701        let p3 = pat("p3", vec!["c".to_string()], 30);
702        let graph = PatternDependencyGraph::build(vec![p1, p2, p3]);
703
704        let total_in_stages: usize = graph.execution_order().iter().map(|s| s.len()).sum();
705        assert_eq!(
706            total_in_stages, 3,
707            "All patterns should appear in execution stages"
708        );
709    }
710
711    #[test]
712    fn test_dependency_graph_patterns_accessor() {
713        let p1 = pat("x", vec!["a".to_string()], 5);
714        let p2 = pat("y", vec!["b".to_string()], 15);
715        let graph = PatternDependencyGraph::build(vec![p1, p2]);
716
717        let patterns = graph.patterns();
718        assert_eq!(patterns.len(), 2);
719        assert_eq!(patterns[0].estimated_cardinality, 5);
720        assert_eq!(patterns[1].estimated_cardinality, 15);
721    }
722
723    // --- merge_results extended tests ---
724
725    #[test]
726    fn test_merge_results_multi_var_join() {
727        let evaluator = ParallelBgpEvaluator::new(1);
728
729        let mut row_l = BindingRow::new();
730        row_l.insert("x".to_string(), RdfTerm::iri("http://a"));
731        row_l.insert("y".to_string(), RdfTerm::iri("http://b"));
732
733        let mut row_r = BindingRow::new();
734        row_r.insert("x".to_string(), RdfTerm::iri("http://a"));
735        row_r.insert("y".to_string(), RdfTerm::iri("http://b"));
736        row_r.insert("z".to_string(), RdfTerm::iri("http://c"));
737
738        let result = evaluator.merge_results(
739            vec![row_l],
740            vec![row_r],
741            &["x".to_string(), "y".to_string()],
742        );
743
744        assert_eq!(
745            result.len(),
746            1,
747            "Matching multi-var join should produce one row"
748        );
749        assert!(
750            result[0].contains_key("z"),
751            "Joined row should contain z from right side"
752        );
753    }
754
755    #[test]
756    fn test_merge_results_no_matching_join_vars() {
757        let evaluator = ParallelBgpEvaluator::new(1);
758
759        let mut row_l = BindingRow::new();
760        row_l.insert("x".to_string(), RdfTerm::iri("http://a"));
761
762        let mut row_r = BindingRow::new();
763        row_r.insert("x".to_string(), RdfTerm::iri("http://DIFFERENT"));
764
765        let result = evaluator.merge_results(vec![row_l], vec![row_r], &["x".to_string()]);
766
767        assert_eq!(
768            result.len(),
769            0,
770            "Non-matching join should produce empty result"
771        );
772    }
773
774    #[test]
775    fn test_merge_results_multiple_right_matches() {
776        let evaluator = ParallelBgpEvaluator::new(1);
777
778        let mut row_l = BindingRow::new();
779        row_l.insert("x".to_string(), RdfTerm::iri("http://shared"));
780
781        let right: Vec<BindingRow> = (0..3)
782            .map(|i| {
783                let mut row = BindingRow::new();
784                row.insert("x".to_string(), RdfTerm::iri("http://shared"));
785                row.insert("y".to_string(), RdfTerm::iri(format!("http://val{i}")));
786                row
787            })
788            .collect();
789
790        let result = evaluator.merge_results(vec![row_l], right, &["x".to_string()]);
791        assert_eq!(
792            result.len(),
793            3,
794            "Should produce one row for each matching right-side row"
795        );
796    }
797
798    // --- ParallelBgpEvaluator configuration tests ---
799
800    #[test]
801    fn test_evaluator_chunk_size_minimum_is_one() {
802        let evaluator = ParallelBgpEvaluator::new(4).with_chunk_size(0);
803        assert_eq!(evaluator.chunk_size, 1, "Chunk size should be at least 1");
804    }
805
806    #[test]
807    fn test_evaluator_chunk_size_set_correctly() {
808        let evaluator = ParallelBgpEvaluator::new(4).with_chunk_size(8);
809        assert_eq!(evaluator.chunk_size, 8);
810    }
811
812    #[test]
813    fn test_evaluator_default_uses_cpu_count() {
814        let evaluator = ParallelBgpEvaluator::default();
815        assert!(evaluator.num_threads >= 1, "Should use at least 1 thread");
816    }
817
818    // --- Evaluate with MockTripleStore ---
819
820    #[test]
821    fn test_evaluate_no_results_from_store() {
822        // MockTripleStore.new() returns empty rows by default for unknown pattern ids.
823        // The evaluator starts with a single empty binding row (identity element for join).
824        // merge_results with right=empty returns left, so we get back the initial empty row.
825        // Verify that the result contains no variable bindings (all rows are empty maps).
826        let store = MockTripleStore::new();
827        let evaluator = ParallelBgpEvaluator::new(1);
828
829        let pattern = pat("no_results", vec!["x".to_string(), "y".to_string()], 100);
830        let result = evaluator.evaluate(vec![pattern], &store).unwrap();
831        // Result may contain the initial empty binding row - verify no actual bindings exist
832        let has_bindings = result.iter().any(|row| !row.is_empty());
833        assert!(
834            !has_bindings,
835            "Empty store should produce no variable bindings"
836        );
837    }
838
839    #[test]
840    fn test_evaluate_two_independent_patterns_cross_product() {
841        let mut store = MockTripleStore::new();
842
843        let p1 = pat("pat_a", vec!["a".to_string()], 2);
844        let p2 = pat("pat_b", vec!["b".to_string()], 3);
845
846        store.results.insert(
847            "pat_a".to_string(),
848            vec![
849                make_row(&[("a", iri_term("http://a1"))]),
850                make_row(&[("a", iri_term("http://a2"))]),
851            ],
852        );
853        store.results.insert(
854            "pat_b".to_string(),
855            vec![
856                make_row(&[("b", iri_term("http://b1"))]),
857                make_row(&[("b", iri_term("http://b2"))]),
858                make_row(&[("b", iri_term("http://b3"))]),
859            ],
860        );
861
862        let evaluator = ParallelBgpEvaluator::new(1);
863        let result = evaluator.evaluate(vec![p1, p2], &store).unwrap();
864
865        // 2 * 3 = 6 cross-product rows (independent patterns)
866        assert_eq!(
867            result.len(),
868            6,
869            "Independent patterns produce cross product"
870        );
871    }
872}