Skip to main content

qualia_core_db/sparql_library/
range_join_select.rs

1//! Range-native nested-loop joins wrapped by Project / Filter / Limit.
2//!
3//! The existing join kernel only accepted a root `NestedLoopJoin`. Real SELECT
4//! trees are `Project(Filter(Limit(Join(Scan, Scan))))`. This module peels those
5//! wrappers without materialising a resident graph.
6
7use super::range_select_apply::{apply_select_wrappers, SelectWrapperState};
8use super::sparql_ast::{BindingRow, SparqlQueryContext, VariableId, MAX_VARIABLES};
9use super::sparql_executor::{
10    execute_range_nested_loop_join_page_into, execute_range_volume_set_nested_loop_join_page_into,
11    Q42RangeNestedLoopJoinPage, Q42RangeNestedLoopJoinPlan, Q42RangeNestedLoopJoinState,
12    Q42RangeVolumeSetNestedLoopJoinState,
13};
14use super::sparql_planner::{ExecutionPlan, PhysicalOperatorType};
15use crate::NQuin;
16
17/// A join plus the SELECT wrappers that sit above it.
18#[derive(Clone, Copy, Debug, Eq, PartialEq)]
19pub struct Q42RangeJoinSelectPlan {
20    pub join: Q42RangeNestedLoopJoinPlan,
21    pub projection: [VariableId; MAX_VARIABLES],
22    pub projection_count: u8,
23    pub filters: [super::sparql_ast::ExpressionId; 8],
24    pub filter_count: u8,
25    pub limit: u64,
26    pub offset: u64,
27}
28
29#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
30pub struct Q42RangeJoinSelectState {
31    pub join: Q42RangeNestedLoopJoinState,
32    pub wrappers: SelectWrapperState,
33}
34
35#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
36pub struct Q42RangeVolumeSetJoinSelectState {
37    pub join: Q42RangeVolumeSetNestedLoopJoinState,
38    pub wrappers: SelectWrapperState,
39}
40
41impl Q42RangeJoinSelectPlan {
42    pub fn from_execution_plan(plan: &ExecutionPlan) -> Result<Self, String> {
43        if plan.operator_count == 0 || plan.root_operator as usize >= plan.operator_count as usize {
44            return Err("range join-select requires a non-empty execution plan".to_string());
45        }
46        let mut operator = plan.root_operator;
47        let mut projection = [0; MAX_VARIABLES];
48        let mut projection_count = 0;
49        let mut filters = [0; 8];
50        let mut filter_count = 0usize;
51        let mut limit = u64::MAX;
52        let mut offset = 0;
53        loop {
54            match plan.operators[operator as usize].operator_type {
55                PhysicalOperatorType::Project {
56                    input,
57                    vars,
58                    var_count,
59                } => {
60                    if projection_count != 0 {
61                        return Err("range join-select does not support nested projections".into());
62                    }
63                    projection = vars;
64                    projection_count = var_count;
65                    operator = input;
66                }
67                PhysicalOperatorType::Limit {
68                    input,
69                    limit: configured,
70                    offset: configured_offset,
71                } => {
72                    if limit != u64::MAX || offset != 0 {
73                        return Err("range join-select does not support nested limits".into());
74                    }
75                    limit = configured;
76                    offset = configured_offset;
77                    operator = input;
78                }
79                PhysicalOperatorType::Filter { input, expression } => {
80                    if filter_count == filters.len() {
81                        return Err(
82                            "range join-select supports at most eight stacked FILTER operators"
83                                .into(),
84                        );
85                    }
86                    filters[filter_count] = expression;
87                    filter_count += 1;
88                    operator = input;
89                }
90                PhysicalOperatorType::NestedLoopJoin { .. } => {
91                    return Ok(Self {
92                        join: Q42RangeNestedLoopJoinPlan::from_join_root(plan, operator)?,
93                        projection,
94                        projection_count,
95                        filters,
96                        filter_count: filter_count as u8,
97                        limit,
98                        offset,
99                    });
100                }
101                PhysicalOperatorType::HashJoin { .. } => {
102                    return Err(
103                        "range join-select does not yet execute HashJoin; planner must emit NestedLoopJoin"
104                            .into(),
105                    );
106                }
107                _ => {
108                    return Err(
109                        "range join-select supports Project/Filter/Limit over NestedLoopJoin(Scan, Scan)"
110                            .into(),
111                    );
112                }
113            }
114        }
115    }
116}
117
118pub fn execute_range_join_select_page_into<S: crate::q42_volume::Q42RangeSource>(
119    volume: &crate::q42_volume::Q42RangeVolume<S>,
120    plan: Q42RangeJoinSelectPlan,
121    ctx: &SparqlQueryContext,
122    state: &mut Q42RangeJoinSelectState,
123    compressed: &mut [u8],
124    decoded: &mut [u8],
125    quin_scratch: &mut [NQuin],
126    left_rows: &mut [BindingRow],
127    right_rows: &mut [BindingRow],
128    join_out: &mut [BindingRow],
129    out: &mut [BindingRow],
130) -> Result<Q42RangeNestedLoopJoinPage, String> {
131    let mut returned = 0usize;
132    loop {
133        if returned == out.len() {
134            return Ok(Q42RangeNestedLoopJoinPage {
135                returned,
136                done: false,
137            });
138        }
139        if state.wrappers.emitted >= plan.limit {
140            return Ok(Q42RangeNestedLoopJoinPage {
141                returned,
142                done: true,
143            });
144        }
145        let page = execute_range_nested_loop_join_page_into(
146            volume,
147            plan.join,
148            ctx,
149            &mut state.join,
150            compressed,
151            decoded,
152            quin_scratch,
153            left_rows,
154            right_rows,
155            join_out,
156        )?;
157        let applied = apply_select_wrappers(
158            ctx,
159            &plan.filters,
160            plan.filter_count,
161            plan.projection,
162            plan.projection_count,
163            plan.limit,
164            plan.offset,
165            &mut state.wrappers,
166            &join_out[..page.returned],
167            &mut out[returned..],
168        )?;
169        returned += applied.returned;
170        if page.done {
171            return Ok(Q42RangeNestedLoopJoinPage {
172                returned,
173                done: true,
174            });
175        }
176        if applied.limit_reached {
177            return Ok(Q42RangeNestedLoopJoinPage {
178                returned,
179                done: true,
180            });
181        }
182        if applied.returned == 0 && !page.done {
183            continue;
184        }
185        if returned > 0 {
186            return Ok(Q42RangeNestedLoopJoinPage {
187                returned,
188                done: false,
189            });
190        }
191    }
192}
193
194pub fn execute_range_volume_set_join_select_page_into<S: crate::q42_volume::Q42RangeSource>(
195    volumes: &crate::q42_volume::Q42RangeVolumeSet<S>,
196    plan: Q42RangeJoinSelectPlan,
197    ctx: &SparqlQueryContext,
198    state: &mut Q42RangeVolumeSetJoinSelectState,
199    compressed: &mut [u8],
200    decoded: &mut [u8],
201    quin_scratch: &mut [NQuin],
202    left_rows: &mut [BindingRow],
203    right_rows: &mut [BindingRow],
204    join_out: &mut [BindingRow],
205    out: &mut [BindingRow],
206) -> Result<Q42RangeNestedLoopJoinPage, String> {
207    let mut returned = 0usize;
208    loop {
209        if returned == out.len() {
210            return Ok(Q42RangeNestedLoopJoinPage {
211                returned,
212                done: false,
213            });
214        }
215        if state.wrappers.emitted >= plan.limit {
216            return Ok(Q42RangeNestedLoopJoinPage {
217                returned,
218                done: true,
219            });
220        }
221        let page = execute_range_volume_set_nested_loop_join_page_into(
222            volumes,
223            plan.join,
224            ctx,
225            &mut state.join,
226            compressed,
227            decoded,
228            quin_scratch,
229            left_rows,
230            right_rows,
231            join_out,
232        )?;
233        let applied = apply_select_wrappers(
234            ctx,
235            &plan.filters,
236            plan.filter_count,
237            plan.projection,
238            plan.projection_count,
239            plan.limit,
240            plan.offset,
241            &mut state.wrappers,
242            &join_out[..page.returned],
243            &mut out[returned..],
244        )?;
245        returned += applied.returned;
246        if page.done || applied.limit_reached {
247            return Ok(Q42RangeNestedLoopJoinPage {
248                returned,
249                done: true,
250            });
251        }
252        if applied.returned == 0 {
253            continue;
254        }
255        if returned > 0 {
256            return Ok(Q42RangeNestedLoopJoinPage {
257                returned,
258                done: false,
259            });
260        }
261    }
262}
263
264#[cfg(test)]
265mod tests {
266    use super::*;
267    use crate::q42_volume::{LocalFileRangeSource, Q42RangeVolume};
268    use crate::sparql_planner::{ExecutionPlan, PhysicalOperatorType};
269
270    #[test]
271    fn peels_project_filter_limit_from_a_join_tree() {
272        let mut plan = ExecutionPlan::new();
273        let left = plan
274            .add_operator(
275                PhysicalOperatorType::TripleScan {
276                    subject: 10,
277                    predicate: 20,
278                    object: 0,
279                },
280                1,
281            )
282            .unwrap();
283        let right = plan
284            .add_operator(
285                PhysicalOperatorType::TripleScan {
286                    subject: 0,
287                    predicate: 40,
288                    object: 50,
289                },
290                1,
291            )
292            .unwrap();
293        let join = plan
294            .add_operator(
295                PhysicalOperatorType::NestedLoopJoin {
296                    left,
297                    right,
298                    join_var: 0,
299                },
300                1,
301            )
302            .unwrap();
303        let limited = plan
304            .add_operator(
305                PhysicalOperatorType::Limit {
306                    input: join,
307                    limit: 1,
308                    offset: 0,
309                },
310                1,
311            )
312            .unwrap();
313        let projected = plan
314            .add_operator(
315                PhysicalOperatorType::Project {
316                    input: limited,
317                    vars: {
318                        let mut vars = [0; MAX_VARIABLES];
319                        vars[0] = 0;
320                        vars
321                    },
322                    var_count: 1,
323                },
324                1,
325            )
326            .unwrap();
327        plan.root_operator = projected;
328        let compiled = Q42RangeJoinSelectPlan::from_execution_plan(&plan).unwrap();
329        assert_eq!(compiled.limit, 1);
330        assert_eq!(compiled.projection_count, 1);
331        assert_eq!(compiled.join.left.subject, 10);
332        assert_eq!(compiled.join.right.object, 50);
333    }
334
335    #[test]
336    fn limit_stops_after_one_join_row() {
337        let dir = tempfile::TempDir::new().unwrap();
338        let path = dir.path().join("join-select.q42");
339        let left = NQuin {
340            subject: 10,
341            predicate: 20,
342            object: 30,
343            context: 0,
344            metadata: 0,
345            parity: 0,
346        };
347        let right = NQuin {
348            subject: 30,
349            predicate: 40,
350            object: 50,
351            context: 0,
352            metadata: 0,
353            parity: 0,
354        };
355        crate::q42_volume::write_unified_volume(
356            &path,
357            &std::collections::HashMap::new(),
358            &[(left.object, right.object)],
359            &[vec![left, right]],
360        )
361        .unwrap();
362        let source = LocalFileRangeSource::open(&path).unwrap();
363        let volume = Q42RangeVolume::open(source).unwrap();
364        let mut context = SparqlQueryContext::new();
365        context.variable_count = 1;
366        let plan = Q42RangeJoinSelectPlan {
367            join: Q42RangeNestedLoopJoinPlan {
368                left: super::super::sparql_executor::Q42RangeTriplePattern {
369                    subject: left.subject,
370                    predicate: left.predicate,
371                    object: 0,
372                },
373                right: super::super::sparql_executor::Q42RangeTriplePattern {
374                    subject: 0,
375                    predicate: right.predicate,
376                    object: right.object,
377                },
378            },
379            projection: {
380                let mut vars = [0; MAX_VARIABLES];
381                vars[0] = 0;
382                vars
383            },
384            projection_count: 1,
385            filters: [0; 8],
386            filter_count: 0,
387            limit: 1,
388            offset: 0,
389        };
390        let mut compressed = [0u8; crate::q42_volume::MAX_COMPRESSED_SUPERBLOCK_SIZE];
391        let mut decoded = [0u8; crate::q42_volume::SUPERBLOCK_SIZE];
392        let mut quins = [NQuin::default(); 1];
393        let mut left_rows = [BindingRow::default(); 1];
394        let mut right_rows = [BindingRow::default(); 1];
395        let mut join_out = [BindingRow::default(); 1];
396        let mut out = [BindingRow::default(); 1];
397        let mut state = Q42RangeJoinSelectState::default();
398        let page = execute_range_join_select_page_into(
399            &volume,
400            plan,
401            &context,
402            &mut state,
403            &mut compressed,
404            &mut decoded,
405            &mut quins,
406            &mut left_rows,
407            &mut right_rows,
408            &mut join_out,
409            &mut out,
410        )
411        .unwrap();
412        assert_eq!(page.returned, 1);
413        assert!(page.done);
414        assert_eq!(out[0].get(0), Some(right.subject));
415    }
416}