1use super::range_select_apply::{apply_select_wrappers, SelectWrapperState};
8use super::sparql_ast::{
9 BindingRow, ExpressionId, SparqlQueryContext, VariableId, MAX_VARIABLES,
10};
11use super::sparql_executor::{
12 execute_range_triple_page_into, execute_range_volume_set_triple_page_into, Q42RangeNestedLoopJoinPage,
13 Q42RangeSparqlCursor, Q42RangeTriplePattern, Q42RangeVolumeSetSparqlCursor,
14};
15use super::sparql_filter::{EvalResult, ExpressionEvaluator};
16use super::sparql_planner::{ExecutionPlan, PhysicalOperatorType};
17use crate::NQuin;
18
19#[derive(Clone, Copy, Debug, Eq, PartialEq)]
20pub struct Q42RangeBindPlan {
21 pub input: Q42RangeTriplePattern,
22 pub var: VariableId,
23 pub expression: ExpressionId,
24 pub projection: [VariableId; MAX_VARIABLES],
25 pub projection_count: u8,
26 pub filters: [ExpressionId; 8],
27 pub filter_count: u8,
28 pub limit: u64,
29 pub offset: u64,
30}
31
32#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
33pub struct Q42RangeBindState {
34 pub scan: Q42RangeSparqlCursor,
35 pub exhausted: bool,
36 pub wrappers: SelectWrapperState,
37}
38
39#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
40pub struct Q42RangeVolumeSetBindState {
41 pub scan: Q42RangeVolumeSetSparqlCursor,
42 pub exhausted: bool,
43 pub wrappers: SelectWrapperState,
44}
45
46impl Q42RangeBindPlan {
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 BIND 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 BIND 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 BIND 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 BIND supports at most eight stacked FILTER operators".into());
87 }
88 filters[filter_count] = expression;
89 filter_count += 1;
90 operator = input;
91 }
92 PhysicalOperatorType::Bind {
93 input,
94 var,
95 expression,
96 } => {
97 return Ok(Self {
98 input: triple_scan(plan, input)?,
99 var,
100 expression,
101 projection,
102 projection_count,
103 filters,
104 filter_count: filter_count as u8,
105 limit,
106 offset,
107 });
108 }
109 _ => {
110 return Err(
111 "range BIND supports Project/Filter/Limit over Bind(TripleScan)".into(),
112 );
113 }
114 }
115 }
116 }
117}
118
119fn triple_scan(plan: &ExecutionPlan, operator: u16) -> Result<Q42RangeTriplePattern, String> {
120 let Some(entry) = plan.operators.get(operator as usize) else {
121 return Err("range BIND input operator is out of bounds".into());
122 };
123 match entry.operator_type {
124 PhysicalOperatorType::TripleScan {
125 subject,
126 predicate,
127 object,
128 } => Ok(Q42RangeTriplePattern {
129 subject,
130 predicate,
131 object,
132 }),
133 _ => Err("range BIND currently requires a TripleScan input".into()),
134 }
135}
136
137pub fn apply_bind(ctx: &SparqlQueryContext, var: VariableId, expression: ExpressionId, row: &mut BindingRow) {
138 match ExpressionEvaluator::evaluate(expression, ctx, row) {
139 Ok(EvalResult::Numeric(n)) | Ok(EvalResult::Iri(n)) | Ok(EvalResult::String(n)) => {
140 row.set(var, n);
141 }
142 Ok(EvalResult::Boolean(b)) => row.set(var, b as u64),
143 Ok(EvalResult::Float(f)) => row.set(var, f.to_bits()),
144 Err(_) => {}
145 }
146}
147
148pub fn execute_range_bind_page_into<S: crate::q42_volume::Q42RangeSource>(
149 volume: &crate::q42_volume::Q42RangeVolume<S>,
150 plan: Q42RangeBindPlan,
151 ctx: &SparqlQueryContext,
152 state: &mut Q42RangeBindState,
153 compressed: &mut [u8],
154 decoded: &mut [u8],
155 quin_scratch: &mut [NQuin],
156 raw: &mut [BindingRow],
157 out: &mut [BindingRow],
158) -> Result<Q42RangeNestedLoopJoinPage, String> {
159 let mut returned = 0usize;
160 loop {
161 if returned == out.len() {
162 return Ok(Q42RangeNestedLoopJoinPage {
163 returned,
164 done: false,
165 });
166 }
167 if state.wrappers.emitted >= plan.limit || state.exhausted {
168 return Ok(Q42RangeNestedLoopJoinPage {
169 returned,
170 done: true,
171 });
172 }
173 let page = execute_range_triple_page_into(
174 volume,
175 plan.input.subject,
176 plan.input.predicate,
177 plan.input.object,
178 None,
179 ctx,
180 &BindingRow::default(),
181 state.scan,
182 compressed,
183 decoded,
184 quin_scratch,
185 raw,
186 )?;
187 state.scan = page.next_cursor.unwrap_or_default();
188 state.exhausted = page.next_cursor.is_none();
189 for row in &mut raw[..page.returned] {
190 apply_bind(ctx, plan.var, plan.expression, row);
191 }
192 let applied = apply_select_wrappers(
193 ctx,
194 &plan.filters,
195 plan.filter_count,
196 plan.projection,
197 plan.projection_count,
198 plan.limit,
199 plan.offset,
200 &mut state.wrappers,
201 &raw[..page.returned],
202 &mut out[returned..],
203 )?;
204 returned += applied.returned;
205 if state.exhausted || applied.limit_reached {
206 return Ok(Q42RangeNestedLoopJoinPage {
207 returned,
208 done: true,
209 });
210 }
211 if applied.returned == 0 {
212 continue;
213 }
214 if returned > 0 {
215 return Ok(Q42RangeNestedLoopJoinPage {
216 returned,
217 done: false,
218 });
219 }
220 }
221}
222
223pub fn execute_range_volume_set_bind_page_into<S: crate::q42_volume::Q42RangeSource>(
224 volumes: &crate::q42_volume::Q42RangeVolumeSet<S>,
225 plan: Q42RangeBindPlan,
226 ctx: &SparqlQueryContext,
227 state: &mut Q42RangeVolumeSetBindState,
228 compressed: &mut [u8],
229 decoded: &mut [u8],
230 quin_scratch: &mut [NQuin],
231 raw: &mut [BindingRow],
232 out: &mut [BindingRow],
233) -> Result<Q42RangeNestedLoopJoinPage, String> {
234 let mut returned = 0usize;
235 loop {
236 if returned == out.len() {
237 return Ok(Q42RangeNestedLoopJoinPage {
238 returned,
239 done: false,
240 });
241 }
242 if state.wrappers.emitted >= plan.limit || state.exhausted {
243 return Ok(Q42RangeNestedLoopJoinPage {
244 returned,
245 done: true,
246 });
247 }
248 let page = execute_range_volume_set_triple_page_into(
249 volumes,
250 plan.input.subject,
251 plan.input.predicate,
252 plan.input.object,
253 None,
254 ctx,
255 &BindingRow::default(),
256 state.scan,
257 compressed,
258 decoded,
259 quin_scratch,
260 raw,
261 )?;
262 state.scan = page.next_cursor.unwrap_or_default();
263 state.exhausted = page.next_cursor.is_none();
264 for row in &mut raw[..page.returned] {
265 apply_bind(ctx, plan.var, plan.expression, row);
266 }
267 let applied = apply_select_wrappers(
268 ctx,
269 &plan.filters,
270 plan.filter_count,
271 plan.projection,
272 plan.projection_count,
273 plan.limit,
274 plan.offset,
275 &mut state.wrappers,
276 &raw[..page.returned],
277 &mut out[returned..],
278 )?;
279 returned += applied.returned;
280 if state.exhausted || applied.limit_reached {
281 return Ok(Q42RangeNestedLoopJoinPage {
282 returned,
283 done: true,
284 });
285 }
286 if applied.returned == 0 {
287 continue;
288 }
289 if returned > 0 {
290 return Ok(Q42RangeNestedLoopJoinPage {
291 returned,
292 done: false,
293 });
294 }
295 }
296}
297
298#[cfg(test)]
299mod tests {
300 use super::*;
301 use crate::q42_volume::{write_unified_volume, LocalFileRangeSource, Q42RangeVolume};
302 use crate::sparql_ast::Expression;
303
304 fn quin(s: u64, p: u64, o: u64) -> NQuin {
305 NQuin {
306 subject: s,
307 predicate: p,
308 object: o,
309 context: 0,
310 metadata: 0,
311 parity: 0,
312 }
313 }
314
315 fn test_volume(quins: &[NQuin]) -> (tempfile::TempDir, Q42RangeVolume<LocalFileRangeSource>) {
316 let dir = tempfile::TempDir::new().unwrap();
317 let path = dir.path().join("bind.q42");
318 let lo = quins.iter().map(|q| q.object).min().unwrap_or(0);
319 let hi = quins.iter().map(|q| q.object).max().unwrap_or(0);
320 write_unified_volume(
321 &path,
322 &std::collections::HashMap::new(),
323 &[(lo, hi)],
324 &[quins.to_vec()],
325 )
326 .unwrap();
327 let source = LocalFileRangeSource::open(&path).unwrap();
328 (dir, Q42RangeVolume::open(source).unwrap())
329 }
330
331 #[test]
332 fn peels_project_from_bind_tree() {
333 let mut plan = ExecutionPlan::new();
334 let scan = plan
335 .add_operator(
336 PhysicalOperatorType::TripleScan {
337 subject: 0,
338 predicate: 20,
339 object: 1,
340 },
341 1,
342 )
343 .unwrap();
344 let bind = plan
345 .add_operator(
346 PhysicalOperatorType::Bind {
347 input: scan,
348 var: 2,
349 expression: 0,
350 },
351 1,
352 )
353 .unwrap();
354 let mut vars = [0; MAX_VARIABLES];
355 vars[0] = 2;
356 plan.root_operator = plan
357 .add_operator(
358 PhysicalOperatorType::Project {
359 input: bind,
360 vars,
361 var_count: 1,
362 },
363 1,
364 )
365 .unwrap();
366 let compiled = Q42RangeBindPlan::from_execution_plan(&plan).unwrap();
367 assert_eq!(compiled.var, 2);
368 assert_eq!(compiled.input.predicate, 20);
369 assert_eq!(compiled.projection_count, 1);
370 }
371
372 #[test]
373 fn bind_literal_extends_every_row() {
374 let row = quin(10, 20, 30);
375 let (_dir, volume) = test_volume(&[row]);
376 let mut ctx = SparqlQueryContext::new();
377 ctx.variable_count = 3;
378 let expr = ctx.alloc_expression(Expression::Literal(42)).unwrap();
379 let plan = Q42RangeBindPlan {
380 input: Q42RangeTriplePattern {
381 subject: 0,
382 predicate: 20,
383 object: 1,
384 },
385 var: 2,
386 expression: expr,
387 projection: [0; MAX_VARIABLES],
388 projection_count: 0,
389 filters: [0; 8],
390 filter_count: 0,
391 limit: u64::MAX,
392 offset: 0,
393 };
394 let mut compressed = [0u8; crate::q42_volume::MAX_COMPRESSED_SUPERBLOCK_SIZE];
395 let mut decoded = [0u8; crate::q42_volume::SUPERBLOCK_SIZE];
396 let mut quins = [NQuin::default(); 2];
397 let mut raw = [BindingRow::default(); 2];
398 let mut out = [BindingRow::default(); 2];
399 let mut state = Q42RangeBindState::default();
400 let page = execute_range_bind_page_into(
401 &volume,
402 plan,
403 &ctx,
404 &mut state,
405 &mut compressed,
406 &mut decoded,
407 &mut quins,
408 &mut raw,
409 &mut out,
410 )
411 .unwrap();
412 assert!(page.done);
413 assert_eq!(page.returned, 1);
414 assert_eq!(out[0].get(0), Some(10));
415 assert_eq!(out[0].get(1), Some(30));
416 assert_eq!(out[0].get(2), Some(42));
417 }
418
419 #[test]
420 fn bind_error_leaves_variable_unbound() {
421 let row = quin(10, 20, 30);
422 let (_dir, volume) = test_volume(&[row]);
423 let mut ctx = SparqlQueryContext::new();
424 ctx.variable_count = 3;
425 let plan = Q42RangeBindPlan {
426 input: Q42RangeTriplePattern {
427 subject: 0,
428 predicate: 20,
429 object: 1,
430 },
431 var: 2,
432 expression: 200,
433 projection: [0; MAX_VARIABLES],
434 projection_count: 0,
435 filters: [0; 8],
436 filter_count: 0,
437 limit: u64::MAX,
438 offset: 0,
439 };
440 let mut compressed = [0u8; crate::q42_volume::MAX_COMPRESSED_SUPERBLOCK_SIZE];
441 let mut decoded = [0u8; crate::q42_volume::SUPERBLOCK_SIZE];
442 let mut quins = [NQuin::default(); 2];
443 let mut raw = [BindingRow::default(); 2];
444 let mut out = [BindingRow::default(); 2];
445 let mut state = Q42RangeBindState::default();
446 let page = execute_range_bind_page_into(
447 &volume,
448 plan,
449 &ctx,
450 &mut state,
451 &mut compressed,
452 &mut decoded,
453 &mut quins,
454 &mut raw,
455 &mut out,
456 )
457 .unwrap();
458 assert_eq!(page.returned, 1);
459 assert_eq!(out[0].get(0), Some(10));
460 assert_eq!(out[0].get(2), None);
461 }
462}