Skip to main content

qualia_core_db/sparql_library/
range_union.rs

1//! Range-native SPARQL UNION over two paged TripleScans.
2//!
3//! SPARQL UNION is a bag union: left pages are emitted first, then right
4//! pages. No deduplication. Project / Filter / Limit wrappers sit above the
5//! union and are applied per page.
6
7use super::range_select_apply::{apply_select_wrappers, SelectWrapperState};
8use super::sparql_ast::{BindingRow, ExpressionId, SparqlQueryContext, VariableId, MAX_VARIABLES};
9use super::sparql_executor::{
10    execute_range_triple_page_into, execute_range_volume_set_triple_page_into, Q42RangeNestedLoopJoinPage,
11    Q42RangeSparqlCursor, Q42RangeTriplePattern, Q42RangeVolumeSetSparqlCursor,
12};
13use super::sparql_planner::{ExecutionPlan, PhysicalOperatorType};
14use crate::NQuin;
15
16#[derive(Clone, Copy, Debug, Eq, PartialEq)]
17pub struct Q42RangeUnionPlan {
18    pub left: Q42RangeTriplePattern,
19    pub right: Q42RangeTriplePattern,
20    pub projection: [VariableId; MAX_VARIABLES],
21    pub projection_count: u8,
22    pub filters: [ExpressionId; 8],
23    pub filter_count: u8,
24    pub limit: u64,
25    pub offset: u64,
26}
27
28#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
29pub struct Q42RangeUnionState {
30    pub left_scan: Q42RangeSparqlCursor,
31    pub right_scan: Q42RangeSparqlCursor,
32    pub left_done: bool,
33    pub right_done: bool,
34    pub wrappers: SelectWrapperState,
35}
36
37#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
38pub struct Q42RangeVolumeSetUnionState {
39    pub left_scan: Q42RangeVolumeSetSparqlCursor,
40    pub right_scan: Q42RangeVolumeSetSparqlCursor,
41    pub left_done: bool,
42    pub right_done: bool,
43    pub wrappers: SelectWrapperState,
44}
45
46impl Q42RangeUnionPlan {
47    pub fn from_execution_plan(plan: &ExecutionPlan) -> Result<Self, String> {
48        if plan.operator_count == 0 || plan.root_operator as usize >= plan.operator_count as usize {
49            return Err("range UNION requires a non-empty execution plan".into());
50        }
51        let mut operator = plan.root_operator;
52        let mut projection = [0; MAX_VARIABLES];
53        let mut projection_count = 0;
54        let mut filters = [0; 8];
55        let mut filter_count = 0usize;
56        let mut limit = u64::MAX;
57        let mut offset = 0;
58        loop {
59            match plan.operators[operator as usize].operator_type {
60                PhysicalOperatorType::Project {
61                    input,
62                    vars,
63                    var_count,
64                } => {
65                    if projection_count != 0 {
66                        return Err("range UNION does not support nested projections".into());
67                    }
68                    projection = vars;
69                    projection_count = var_count;
70                    operator = input;
71                }
72                PhysicalOperatorType::Limit {
73                    input,
74                    limit: configured,
75                    offset: configured_offset,
76                } => {
77                    if limit != u64::MAX || offset != 0 {
78                        return Err("range UNION does not support nested limits".into());
79                    }
80                    limit = configured;
81                    offset = configured_offset;
82                    operator = input;
83                }
84                PhysicalOperatorType::Filter { input, expression } => {
85                    if filter_count == filters.len() {
86                        return Err("range UNION supports at most eight stacked FILTER operators".into());
87                    }
88                    filters[filter_count] = expression;
89                    filter_count += 1;
90                    operator = input;
91                }
92                PhysicalOperatorType::Union { left, right } => {
93                    return Ok(Self {
94                        left: triple_scan(plan, left)?,
95                        right: triple_scan(plan, right)?,
96                        projection,
97                        projection_count,
98                        filters,
99                        filter_count: filter_count as u8,
100                        limit,
101                        offset,
102                    });
103                }
104                _ => {
105                    return Err(
106                        "range UNION supports Project/Filter/Limit over Union(TripleScan, TripleScan)"
107                            .into(),
108                    );
109                }
110            }
111        }
112    }
113}
114
115fn triple_scan(plan: &ExecutionPlan, operator: u16) -> Result<Q42RangeTriplePattern, String> {
116    let Some(entry) = plan.operators.get(operator as usize) else {
117        return Err("range UNION input operator is out of bounds".into());
118    };
119    match entry.operator_type {
120        PhysicalOperatorType::TripleScan {
121            subject,
122            predicate,
123            object,
124        } => Ok(Q42RangeTriplePattern {
125            subject,
126            predicate,
127            object,
128        }),
129        _ => Err("range UNION currently requires TripleScan inputs".into()),
130    }
131}
132
133pub fn execute_range_union_page_into<S: crate::q42_volume::Q42RangeSource>(
134    volume: &crate::q42_volume::Q42RangeVolume<S>,
135    plan: Q42RangeUnionPlan,
136    ctx: &SparqlQueryContext,
137    state: &mut Q42RangeUnionState,
138    compressed: &mut [u8],
139    decoded: &mut [u8],
140    quin_scratch: &mut [NQuin],
141    raw: &mut [BindingRow],
142    out: &mut [BindingRow],
143) -> Result<Q42RangeNestedLoopJoinPage, String> {
144    let mut returned = 0usize;
145    loop {
146        if returned == out.len() {
147            return Ok(Q42RangeNestedLoopJoinPage {
148                returned,
149                done: false,
150            });
151        }
152        if state.wrappers.emitted >= plan.limit {
153            return Ok(Q42RangeNestedLoopJoinPage {
154                returned,
155                done: true,
156            });
157        }
158        let produced = fill_union_page(
159            |pattern, cursor, rows| {
160                execute_range_triple_page_into(
161                    volume,
162                    pattern.subject,
163                    pattern.predicate,
164                    pattern.object,
165                    None,
166                    ctx,
167                    &BindingRow::default(),
168                    cursor,
169                    compressed,
170                    decoded,
171                    quin_scratch,
172                    rows,
173                )
174                .map(|page| (page.returned, page.next_cursor))
175            },
176            plan,
177            &mut state.left_scan,
178            &mut state.right_scan,
179            &mut state.left_done,
180            &mut state.right_done,
181            raw,
182        )?;
183        let applied = apply_select_wrappers(
184            ctx,
185            &plan.filters,
186            plan.filter_count,
187            plan.projection,
188            plan.projection_count,
189            plan.limit,
190            plan.offset,
191            &mut state.wrappers,
192            &raw[..produced.returned],
193            &mut out[returned..],
194        )?;
195        returned += applied.returned;
196        if produced.done || applied.limit_reached {
197            return Ok(Q42RangeNestedLoopJoinPage {
198                returned,
199                done: true,
200            });
201        }
202        if applied.returned == 0 {
203            continue;
204        }
205        if returned > 0 {
206            return Ok(Q42RangeNestedLoopJoinPage {
207                returned,
208                done: false,
209            });
210        }
211    }
212}
213
214pub fn execute_range_volume_set_union_page_into<S: crate::q42_volume::Q42RangeSource>(
215    volumes: &crate::q42_volume::Q42RangeVolumeSet<S>,
216    plan: Q42RangeUnionPlan,
217    ctx: &SparqlQueryContext,
218    state: &mut Q42RangeVolumeSetUnionState,
219    compressed: &mut [u8],
220    decoded: &mut [u8],
221    quin_scratch: &mut [NQuin],
222    raw: &mut [BindingRow],
223    out: &mut [BindingRow],
224) -> Result<Q42RangeNestedLoopJoinPage, String> {
225    let mut returned = 0usize;
226    loop {
227        if returned == out.len() {
228            return Ok(Q42RangeNestedLoopJoinPage {
229                returned,
230                done: false,
231            });
232        }
233        if state.wrappers.emitted >= plan.limit {
234            return Ok(Q42RangeNestedLoopJoinPage {
235                returned,
236                done: true,
237            });
238        }
239        let produced = fill_union_page(
240            |pattern, cursor, rows| {
241                execute_range_volume_set_triple_page_into(
242                    volumes,
243                    pattern.subject,
244                    pattern.predicate,
245                    pattern.object,
246                    None,
247                    ctx,
248                    &BindingRow::default(),
249                    cursor,
250                    compressed,
251                    decoded,
252                    quin_scratch,
253                    rows,
254                )
255                .map(|page| (page.returned, page.next_cursor))
256            },
257            plan,
258            &mut state.left_scan,
259            &mut state.right_scan,
260            &mut state.left_done,
261            &mut state.right_done,
262            raw,
263        )?;
264        let applied = apply_select_wrappers(
265            ctx,
266            &plan.filters,
267            plan.filter_count,
268            plan.projection,
269            plan.projection_count,
270            plan.limit,
271            plan.offset,
272            &mut state.wrappers,
273            &raw[..produced.returned],
274            &mut out[returned..],
275        )?;
276        returned += applied.returned;
277        if produced.done || applied.limit_reached {
278            return Ok(Q42RangeNestedLoopJoinPage {
279                returned,
280                done: true,
281            });
282        }
283        if applied.returned == 0 {
284            continue;
285        }
286        if returned > 0 {
287            return Ok(Q42RangeNestedLoopJoinPage {
288                returned,
289                done: false,
290            });
291        }
292    }
293}
294
295fn fill_union_page<C, F>(
296    mut page_fn: F,
297    plan: Q42RangeUnionPlan,
298    left_scan: &mut C,
299    right_scan: &mut C,
300    left_done: &mut bool,
301    right_done: &mut bool,
302    out: &mut [BindingRow],
303) -> Result<Q42RangeNestedLoopJoinPage, String>
304where
305    C: Copy + Default,
306    F: FnMut(Q42RangeTriplePattern, C, &mut [BindingRow]) -> Result<(usize, Option<C>), String>,
307{
308    if !*left_done {
309        let (count, next) = page_fn(plan.left, *left_scan, out)?;
310        *left_scan = next.unwrap_or_default();
311        *left_done = next.is_none();
312        if count > 0 || !*left_done {
313            return Ok(Q42RangeNestedLoopJoinPage {
314                returned: count,
315                done: false,
316            });
317        }
318    }
319    if !*right_done {
320        let (count, next) = page_fn(plan.right, *right_scan, out)?;
321        *right_scan = next.unwrap_or_default();
322        *right_done = next.is_none();
323        return Ok(Q42RangeNestedLoopJoinPage {
324            returned: count,
325            done: *right_done && count == 0,
326        });
327    }
328    Ok(Q42RangeNestedLoopJoinPage {
329        returned: 0,
330        done: true,
331    })
332}
333
334#[cfg(test)]
335mod tests {
336    use super::*;
337    use crate::q42_volume::{write_unified_volume, LocalFileRangeSource, Q42RangeVolume};
338
339    fn quin(s: u64, p: u64, o: u64) -> NQuin {
340        NQuin {
341            subject: s,
342            predicate: p,
343            object: o,
344            context: 0,
345            metadata: 0,
346            parity: 0,
347        }
348    }
349
350    fn test_volume(quins: &[NQuin]) -> (tempfile::TempDir, Q42RangeVolume<LocalFileRangeSource>) {
351        let dir = tempfile::TempDir::new().unwrap();
352        let path = dir.path().join("union.q42");
353        let lo = quins.iter().map(|q| q.object).min().unwrap_or(0);
354        let hi = quins.iter().map(|q| q.object).max().unwrap_or(0);
355        write_unified_volume(
356            &path,
357            &std::collections::HashMap::new(),
358            &[(lo, hi)],
359            &[quins.to_vec()],
360        )
361        .unwrap();
362        let source = LocalFileRangeSource::open(&path).unwrap();
363        (dir, Q42RangeVolume::open(source).unwrap())
364    }
365
366    #[test]
367    fn peels_project_from_union_tree() {
368        let mut plan = ExecutionPlan::new();
369        let left = plan
370            .add_operator(
371                PhysicalOperatorType::TripleScan {
372                    subject: 0,
373                    predicate: 20,
374                    object: 1,
375                },
376                1,
377            )
378            .unwrap();
379        let right = plan
380            .add_operator(
381                PhysicalOperatorType::TripleScan {
382                    subject: 0,
383                    predicate: 40,
384                    object: 1,
385                },
386                1,
387            )
388            .unwrap();
389        let union = plan
390            .add_operator(PhysicalOperatorType::Union { left, right }, 2)
391            .unwrap();
392        let mut vars = [0; MAX_VARIABLES];
393        vars[0] = 0;
394        plan.root_operator = plan
395            .add_operator(
396                PhysicalOperatorType::Project {
397                    input: union,
398                    vars,
399                    var_count: 1,
400                },
401                2,
402            )
403            .unwrap();
404        let compiled = Q42RangeUnionPlan::from_execution_plan(&plan).unwrap();
405        assert_eq!(compiled.left.predicate, 20);
406        assert_eq!(compiled.right.predicate, 40);
407        assert_eq!(compiled.projection_count, 1);
408    }
409
410    #[test]
411    fn concatenates_left_then_right_without_dedup() {
412        let left = quin(10, 20, 30);
413        let right = quin(11, 40, 31);
414        let (_dir, volume) = test_volume(&[left, right]);
415        let mut ctx = SparqlQueryContext::new();
416        ctx.variable_count = 2;
417        let plan = Q42RangeUnionPlan {
418            left: Q42RangeTriplePattern {
419                subject: 0,
420                predicate: 20,
421                object: 1,
422            },
423            right: Q42RangeTriplePattern {
424                subject: 0,
425                predicate: 40,
426                object: 1,
427            },
428            projection: [0; MAX_VARIABLES],
429            projection_count: 0,
430            filters: [0; 8],
431            filter_count: 0,
432            limit: u64::MAX,
433            offset: 0,
434        };
435        let mut compressed = [0u8; crate::q42_volume::MAX_COMPRESSED_SUPERBLOCK_SIZE];
436        let mut decoded = [0u8; crate::q42_volume::SUPERBLOCK_SIZE];
437        let mut quins = [NQuin::default(); 4];
438        let mut raw = [BindingRow::default(); 4];
439        let mut out = [BindingRow::default(); 4];
440        let mut state = Q42RangeUnionState::default();
441        let mut rows = Vec::new();
442        loop {
443            let page = execute_range_union_page_into(
444                &volume,
445                plan,
446                &ctx,
447                &mut state,
448                &mut compressed,
449                &mut decoded,
450                &mut quins,
451                &mut raw,
452                &mut out,
453            )
454            .unwrap();
455            rows.extend_from_slice(&out[..page.returned]);
456            if page.done {
457                break;
458            }
459        }
460        assert_eq!(rows.len(), 2);
461        assert_eq!(rows[0].get(0), Some(10));
462        assert_eq!(rows[1].get(0), Some(11));
463    }
464
465    #[test]
466    fn empty_both_sides_is_done() {
467        let only = quin(10, 99, 30);
468        let (_dir, volume) = test_volume(&[only]);
469        let mut ctx = SparqlQueryContext::new();
470        ctx.variable_count = 1;
471        let plan = Q42RangeUnionPlan {
472            left: Q42RangeTriplePattern {
473                subject: 0,
474                predicate: 20,
475                object: 1,
476            },
477            right: Q42RangeTriplePattern {
478                subject: 0,
479                predicate: 40,
480                object: 1,
481            },
482            projection: [0; MAX_VARIABLES],
483            projection_count: 0,
484            filters: [0; 8],
485            filter_count: 0,
486            limit: u64::MAX,
487            offset: 0,
488        };
489        let mut compressed = [0u8; crate::q42_volume::MAX_COMPRESSED_SUPERBLOCK_SIZE];
490        let mut decoded = [0u8; crate::q42_volume::SUPERBLOCK_SIZE];
491        let mut quins = [NQuin::default(); 2];
492        let mut raw = [BindingRow::default(); 2];
493        let mut out = [BindingRow::default(); 2];
494        let mut state = Q42RangeUnionState::default();
495        let page = execute_range_union_page_into(
496            &volume,
497            plan,
498            &ctx,
499            &mut state,
500            &mut compressed,
501            &mut decoded,
502            &mut quins,
503            &mut raw,
504            &mut out,
505        )
506        .unwrap();
507        assert!(page.done);
508        assert_eq!(page.returned, 0);
509    }
510}