Skip to main content

qualia_core_db/sparql_library/
range_optional.rs

1//! Range-native SPARQL OPTIONAL (left join) over paged TripleScans.
2//!
3//! Each left row is emitted even when the right pattern produces no match.
4//! Matching uses the existing range triple kernel so already-bound variables
5//! stay compatible. Project / Filter / Limit wrappers are peeled, not
6//! materialised.
7
8use super::range_select_apply::{apply_select_wrappers, SelectWrapperState};
9use super::sparql_ast::{BindingRow, ExpressionId, SparqlQueryContext, VariableId, MAX_VARIABLES};
10use super::sparql_executor::{
11    execute_range_triple_page_into, execute_range_volume_set_triple_page_into, Q42RangeNestedLoopJoinPage,
12    Q42RangeSparqlCursor, Q42RangeTriplePattern, Q42RangeVolumeSetSparqlCursor,
13};
14use super::sparql_planner::{ExecutionPlan, PhysicalOperatorType};
15use crate::NQuin;
16
17#[derive(Clone, Copy, Debug, Eq, PartialEq)]
18pub struct Q42RangeOptionalPlan {
19    pub left: Q42RangeTriplePattern,
20    pub right: Q42RangeTriplePattern,
21    pub projection: [VariableId; MAX_VARIABLES],
22    pub projection_count: u8,
23    pub filters: [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 Q42RangeOptionalState {
31    pub left_scan: Q42RangeSparqlCursor,
32    pub right_scan: Q42RangeSparqlCursor,
33    pub left_count: usize,
34    pub left_index: usize,
35    pub left_exhausted: bool,
36    pub right_active: bool,
37    pub left_row_matched: bool,
38    pub wrappers: SelectWrapperState,
39}
40
41#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
42pub struct Q42RangeVolumeSetOptionalState {
43    pub left_scan: Q42RangeVolumeSetSparqlCursor,
44    pub right_scan: Q42RangeVolumeSetSparqlCursor,
45    pub left_count: usize,
46    pub left_index: usize,
47    pub left_exhausted: bool,
48    pub right_active: bool,
49    pub left_row_matched: bool,
50    pub wrappers: SelectWrapperState,
51}
52
53impl Q42RangeOptionalPlan {
54    pub fn from_execution_plan(plan: &ExecutionPlan) -> Result<Self, String> {
55        if plan.operator_count == 0 || plan.root_operator as usize >= plan.operator_count as usize {
56            return Err("range OPTIONAL requires a non-empty execution plan".into());
57        }
58        let mut operator = plan.root_operator;
59        let mut projection = [0; MAX_VARIABLES];
60        let mut projection_count = 0;
61        let mut filters = [0; 8];
62        let mut filter_count = 0usize;
63        let mut limit = u64::MAX;
64        let mut offset = 0;
65        loop {
66            match plan.operators[operator as usize].operator_type {
67                PhysicalOperatorType::Project {
68                    input,
69                    vars,
70                    var_count,
71                } => {
72                    if projection_count != 0 {
73                        return Err("range OPTIONAL does not support nested projections".into());
74                    }
75                    projection = vars;
76                    projection_count = var_count;
77                    operator = input;
78                }
79                PhysicalOperatorType::Limit {
80                    input,
81                    limit: configured,
82                    offset: configured_offset,
83                } => {
84                    if limit != u64::MAX || offset != 0 {
85                        return Err("range OPTIONAL does not support nested limits".into());
86                    }
87                    limit = configured;
88                    offset = configured_offset;
89                    operator = input;
90                }
91                PhysicalOperatorType::Filter { input, expression } => {
92                    if filter_count == filters.len() {
93                        return Err("range OPTIONAL supports at most eight stacked FILTER operators".into());
94                    }
95                    filters[filter_count] = expression;
96                    filter_count += 1;
97                    operator = input;
98                }
99                PhysicalOperatorType::Optional { left, right } => {
100                    return Ok(Self {
101                        left: triple_scan(plan, left)?,
102                        right: triple_scan(plan, right)?,
103                        projection,
104                        projection_count,
105                        filters,
106                        filter_count: filter_count as u8,
107                        limit,
108                        offset,
109                    });
110                }
111                _ => {
112                    return Err(
113                        "range OPTIONAL supports Project/Filter/Limit over Optional(TripleScan, TripleScan)"
114                            .into(),
115                    );
116                }
117            }
118        }
119    }
120}
121
122fn triple_scan(plan: &ExecutionPlan, operator: u16) -> Result<Q42RangeTriplePattern, String> {
123    let Some(entry) = plan.operators.get(operator as usize) else {
124        return Err("range OPTIONAL input operator is out of bounds".into());
125    };
126    match entry.operator_type {
127        PhysicalOperatorType::TripleScan {
128            subject,
129            predicate,
130            object,
131        } => Ok(Q42RangeTriplePattern {
132            subject,
133            predicate,
134            object,
135        }),
136        _ => Err("range OPTIONAL currently requires TripleScan inputs".into()),
137    }
138}
139
140pub fn execute_range_optional_page_into<S: crate::q42_volume::Q42RangeSource>(
141    volume: &crate::q42_volume::Q42RangeVolume<S>,
142    plan: Q42RangeOptionalPlan,
143    ctx: &SparqlQueryContext,
144    state: &mut Q42RangeOptionalState,
145    compressed: &mut [u8],
146    decoded: &mut [u8],
147    quin_scratch: &mut [NQuin],
148    left_rows: &mut [BindingRow],
149    right_rows: &mut [BindingRow],
150    join_out: &mut [BindingRow],
151    out: &mut [BindingRow],
152) -> Result<Q42RangeNestedLoopJoinPage, String> {
153    if quin_scratch.is_empty()
154        || left_rows.len() < quin_scratch.len()
155        || right_rows.len() < quin_scratch.len()
156        || join_out.len() < quin_scratch.len()
157    {
158        return Err(
159            "range OPTIONAL requires row buffers at least as large as Quin scratch".into(),
160        );
161    }
162    let mut returned = 0usize;
163    loop {
164        if returned == out.len() {
165            return Ok(Q42RangeNestedLoopJoinPage {
166                returned,
167                done: false,
168            });
169        }
170        if state.wrappers.emitted >= plan.limit {
171            return Ok(Q42RangeNestedLoopJoinPage {
172                returned,
173                done: true,
174            });
175        }
176        let scratch_len = quin_scratch.len();
177        let produced = fill_optional_page(
178            |pattern, input, cursor, rows| {
179                execute_range_triple_page_into(
180                    volume,
181                    pattern.subject,
182                    pattern.predicate,
183                    pattern.object,
184                    None,
185                    ctx,
186                    input,
187                    cursor,
188                    compressed,
189                    decoded,
190                    quin_scratch,
191                    rows,
192                )
193                .map(|page| (page.returned, page.next_cursor))
194            },
195            plan,
196            scratch_len,
197            &mut state.left_scan,
198            &mut state.right_scan,
199            &mut state.left_count,
200            &mut state.left_index,
201            &mut state.left_exhausted,
202            &mut state.right_active,
203            &mut state.left_row_matched,
204            left_rows,
205            right_rows,
206            join_out,
207        )?;
208        let applied = apply_select_wrappers(
209            ctx,
210            &plan.filters,
211            plan.filter_count,
212            plan.projection,
213            plan.projection_count,
214            plan.limit,
215            plan.offset,
216            &mut state.wrappers,
217            &join_out[..produced.returned],
218            &mut out[returned..],
219        )?;
220        returned += applied.returned;
221        if produced.done || applied.limit_reached {
222            return Ok(Q42RangeNestedLoopJoinPage {
223                returned,
224                done: true,
225            });
226        }
227        if applied.returned == 0 {
228            continue;
229        }
230        if returned > 0 {
231            return Ok(Q42RangeNestedLoopJoinPage {
232                returned,
233                done: false,
234            });
235        }
236    }
237}
238
239pub fn execute_range_volume_set_optional_page_into<S: crate::q42_volume::Q42RangeSource>(
240    volumes: &crate::q42_volume::Q42RangeVolumeSet<S>,
241    plan: Q42RangeOptionalPlan,
242    ctx: &SparqlQueryContext,
243    state: &mut Q42RangeVolumeSetOptionalState,
244    compressed: &mut [u8],
245    decoded: &mut [u8],
246    quin_scratch: &mut [NQuin],
247    left_rows: &mut [BindingRow],
248    right_rows: &mut [BindingRow],
249    join_out: &mut [BindingRow],
250    out: &mut [BindingRow],
251) -> Result<Q42RangeNestedLoopJoinPage, String> {
252    if quin_scratch.is_empty()
253        || left_rows.len() < quin_scratch.len()
254        || right_rows.len() < quin_scratch.len()
255        || join_out.len() < quin_scratch.len()
256    {
257        return Err(
258            "range OPTIONAL requires row buffers at least as large as Quin scratch".into(),
259        );
260    }
261    let mut returned = 0usize;
262    loop {
263        if returned == out.len() {
264            return Ok(Q42RangeNestedLoopJoinPage {
265                returned,
266                done: false,
267            });
268        }
269        if state.wrappers.emitted >= plan.limit {
270            return Ok(Q42RangeNestedLoopJoinPage {
271                returned,
272                done: true,
273            });
274        }
275        let scratch_len = quin_scratch.len();
276        let produced = fill_optional_page(
277            |pattern, input, cursor, rows| {
278                execute_range_volume_set_triple_page_into(
279                    volumes,
280                    pattern.subject,
281                    pattern.predicate,
282                    pattern.object,
283                    None,
284                    ctx,
285                    input,
286                    cursor,
287                    compressed,
288                    decoded,
289                    quin_scratch,
290                    rows,
291                )
292                .map(|page| (page.returned, page.next_cursor))
293            },
294            plan,
295            scratch_len,
296            &mut state.left_scan,
297            &mut state.right_scan,
298            &mut state.left_count,
299            &mut state.left_index,
300            &mut state.left_exhausted,
301            &mut state.right_active,
302            &mut state.left_row_matched,
303            left_rows,
304            right_rows,
305            join_out,
306        )?;
307        let applied = apply_select_wrappers(
308            ctx,
309            &plan.filters,
310            plan.filter_count,
311            plan.projection,
312            plan.projection_count,
313            plan.limit,
314            plan.offset,
315            &mut state.wrappers,
316            &join_out[..produced.returned],
317            &mut out[returned..],
318        )?;
319        returned += applied.returned;
320        if produced.done || applied.limit_reached {
321            return Ok(Q42RangeNestedLoopJoinPage {
322                returned,
323                done: true,
324            });
325        }
326        if applied.returned == 0 {
327            continue;
328        }
329        if returned > 0 {
330            return Ok(Q42RangeNestedLoopJoinPage {
331                returned,
332                done: false,
333            });
334        }
335    }
336}
337
338fn fill_optional_page<C, F>(
339    mut page_fn: F,
340    plan: Q42RangeOptionalPlan,
341    scratch_len: usize,
342    left_scan: &mut C,
343    right_scan: &mut C,
344    left_count: &mut usize,
345    left_index: &mut usize,
346    left_exhausted: &mut bool,
347    right_active: &mut bool,
348    left_row_matched: &mut bool,
349    left_rows: &mut [BindingRow],
350    right_rows: &mut [BindingRow],
351    out: &mut [BindingRow],
352) -> Result<Q42RangeNestedLoopJoinPage, String>
353where
354    C: Copy + Default,
355    F: FnMut(
356        Q42RangeTriplePattern,
357        &BindingRow,
358        C,
359        &mut [BindingRow],
360    ) -> Result<(usize, Option<C>), String>,
361{
362    let mut returned = 0usize;
363    loop {
364        if returned + scratch_len > out.len() {
365            return Ok(Q42RangeNestedLoopJoinPage {
366                returned,
367                done: false,
368            });
369        }
370        if *left_index >= *left_count {
371            if *left_exhausted {
372                return Ok(Q42RangeNestedLoopJoinPage {
373                    returned,
374                    done: true,
375                });
376            }
377            let (count, next) = page_fn(
378                plan.left,
379                &BindingRow::default(),
380                *left_scan,
381                left_rows,
382            )?;
383            *left_count = count;
384            *left_index = 0;
385            *left_scan = next.unwrap_or_default();
386            *left_exhausted = next.is_none();
387            *right_active = false;
388            *left_row_matched = false;
389            if *left_count == 0 {
390                if *left_exhausted {
391                    return Ok(Q42RangeNestedLoopJoinPage {
392                        returned,
393                        done: true,
394                    });
395                }
396                continue;
397            }
398        }
399        let input = left_rows[*left_index];
400        if !*right_active {
401            *left_row_matched = false;
402        }
403        let (count, next) = page_fn(
404            plan.right,
405            &input,
406            if *right_active {
407                *right_scan
408            } else {
409                C::default()
410            },
411            right_rows,
412        )?;
413        out[returned..returned + count].copy_from_slice(&right_rows[..count]);
414        returned += count;
415        if count > 0 {
416            *left_row_matched = true;
417        }
418        match next {
419            Some(cursor) => {
420                *right_scan = cursor;
421                *right_active = true;
422            }
423            None => {
424                if !*left_row_matched {
425                    if returned == out.len() {
426                        return Ok(Q42RangeNestedLoopJoinPage {
427                            returned,
428                            done: false,
429                        });
430                    }
431                    out[returned] = input;
432                    returned += 1;
433                }
434                *left_index += 1;
435                *right_scan = C::default();
436                *right_active = false;
437                *left_row_matched = false;
438            }
439        }
440        if returned + scratch_len > out.len() {
441            return Ok(Q42RangeNestedLoopJoinPage {
442                returned,
443                done: false,
444            });
445        }
446    }
447}
448
449#[cfg(test)]
450mod tests {
451    use super::*;
452    use crate::q42_volume::{write_unified_volume, LocalFileRangeSource, Q42RangeVolume};
453
454    fn test_volume(quins: &[NQuin]) -> (tempfile::TempDir, Q42RangeVolume<LocalFileRangeSource>) {
455        let dir = tempfile::TempDir::new().unwrap();
456        let path = dir.path().join("optional.q42");
457        let lo = quins.iter().map(|q| q.object).min().unwrap_or(0);
458        let hi = quins.iter().map(|q| q.object).max().unwrap_or(0);
459        write_unified_volume(
460            &path,
461            &std::collections::HashMap::new(),
462            &[(lo, hi)],
463            &[quins.to_vec()],
464        )
465        .unwrap();
466        let source = LocalFileRangeSource::open(&path).unwrap();
467        let volume = Q42RangeVolume::open(source).unwrap();
468        (dir, volume)
469    }
470
471    fn quin(s: u64, p: u64, o: u64) -> NQuin {
472        NQuin {
473            subject: s,
474            predicate: p,
475            object: o,
476            context: 0,
477            metadata: 0,
478            parity: 0,
479        }
480    }
481
482    #[test]
483    fn peels_project_from_optional_tree() {
484        let mut plan = ExecutionPlan::new();
485        let left = plan
486            .add_operator(
487                PhysicalOperatorType::TripleScan {
488                    subject: 0,
489                    predicate: 20,
490                    object: 1,
491                },
492                1,
493            )
494            .unwrap();
495        let right = plan
496            .add_operator(
497                PhysicalOperatorType::TripleScan {
498                    subject: 1,
499                    predicate: 40,
500                    object: 2,
501                },
502                1,
503            )
504            .unwrap();
505        let optional = plan
506            .add_operator(PhysicalOperatorType::Optional { left, right }, 1)
507            .unwrap();
508        let mut vars = [0; MAX_VARIABLES];
509        vars[0] = 0;
510        plan.root_operator = plan
511            .add_operator(
512                PhysicalOperatorType::Project {
513                    input: optional,
514                    vars,
515                    var_count: 1,
516                },
517                1,
518            )
519            .unwrap();
520        let compiled = Q42RangeOptionalPlan::from_execution_plan(&plan).unwrap();
521        assert_eq!(compiled.projection_count, 1);
522        assert_eq!(compiled.left.predicate, 20);
523        assert_eq!(compiled.right.predicate, 40);
524    }
525
526    #[test]
527    fn unmatched_left_row_is_kept() {
528        let required = quin(10, 20, 30);
529        let unmatched = quin(11, 20, 31);
530        let optional = quin(30, 40, 50);
531        let (_dir, volume) = test_volume(&[required, unmatched, optional]);
532        let mut ctx = SparqlQueryContext::new();
533        ctx.variable_count = 2;
534        let plan = Q42RangeOptionalPlan {
535            left: Q42RangeTriplePattern {
536                subject: 0,
537                predicate: 20,
538                object: 1,
539            },
540            right: Q42RangeTriplePattern {
541                subject: 1,
542                predicate: 40,
543                object: 50,
544            },
545            projection: [0; MAX_VARIABLES],
546            projection_count: 0,
547            filters: [0; 8],
548            filter_count: 0,
549            limit: u64::MAX,
550            offset: 0,
551        };
552        let mut compressed = [0u8; crate::q42_volume::MAX_COMPRESSED_SUPERBLOCK_SIZE];
553        let mut decoded = [0u8; crate::q42_volume::SUPERBLOCK_SIZE];
554        let mut quins = [NQuin::default(); 4];
555        let mut left_rows = [BindingRow::default(); 4];
556        let mut right_rows = [BindingRow::default(); 4];
557        let mut join_out = [BindingRow::default(); 4];
558        let mut out = [BindingRow::default(); 4];
559        let mut state = Q42RangeOptionalState::default();
560        let mut rows = Vec::new();
561        loop {
562            let page = execute_range_optional_page_into(
563                &volume,
564                plan,
565                &ctx,
566                &mut state,
567                &mut compressed,
568                &mut decoded,
569                &mut quins,
570                &mut left_rows,
571                &mut right_rows,
572                &mut join_out,
573                &mut out,
574            )
575            .unwrap();
576            rows.extend_from_slice(&out[..page.returned]);
577            if page.done {
578                break;
579            }
580        }
581        assert_eq!(rows.len(), 2);
582        let subjects: Vec<u64> = rows.iter().filter_map(|row| row.get(0)).collect();
583        assert!(subjects.contains(&10));
584        assert!(subjects.contains(&11));
585        let matched = rows.iter().find(|row| row.get(0) == Some(10)).unwrap();
586        assert_eq!(matched.get(1), Some(30));
587        let kept = rows.iter().find(|row| row.get(0) == Some(11)).unwrap();
588        assert_eq!(kept.get(1), Some(31));
589    }
590
591    #[test]
592    fn empty_left_yields_no_rows() {
593        let only_right = quin(30, 40, 50);
594        let (_dir, volume) = test_volume(&[only_right]);
595        let mut ctx = SparqlQueryContext::new();
596        ctx.variable_count = 1;
597        let plan = Q42RangeOptionalPlan {
598            left: Q42RangeTriplePattern {
599                subject: 0,
600                predicate: 20,
601                object: 99,
602            },
603            right: Q42RangeTriplePattern {
604                subject: 0,
605                predicate: 40,
606                object: 50,
607            },
608            projection: [0; MAX_VARIABLES],
609            projection_count: 0,
610            filters: [0; 8],
611            filter_count: 0,
612            limit: u64::MAX,
613            offset: 0,
614        };
615        let mut compressed = [0u8; crate::q42_volume::MAX_COMPRESSED_SUPERBLOCK_SIZE];
616        let mut decoded = [0u8; crate::q42_volume::SUPERBLOCK_SIZE];
617        let mut quins = [NQuin::default(); 2];
618        let mut left_rows = [BindingRow::default(); 2];
619        let mut right_rows = [BindingRow::default(); 2];
620        let mut join_out = [BindingRow::default(); 2];
621        let mut out = [BindingRow::default(); 2];
622        let mut state = Q42RangeOptionalState::default();
623        let page = execute_range_optional_page_into(
624            &volume,
625            plan,
626            &ctx,
627            &mut state,
628            &mut compressed,
629            &mut decoded,
630            &mut quins,
631            &mut left_rows,
632            &mut right_rows,
633            &mut join_out,
634            &mut out,
635        )
636        .unwrap();
637        assert!(page.done);
638        assert_eq!(page.returned, 0);
639    }
640}