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