1use 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}