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