1use crate::compiler::{Col, Compiler, QueryOptions};
10use crate::encoding::{
11 graph_id, named_node_id, EncodedRows, Tag, DEFAULT_GRAPH_ID, PAYLOAD_BITS, PAYLOAD_MASK,
12};
13use crate::error::{Error, Result};
14use crate::sql::{Capabilities, Statement};
15use crate::stats::Stats;
16use crate::writer::{term_statements, EncodedQuads};
17use oxrdf::{BlankNode, GraphName, NamedOrBlankNode, Quad, Term, Variable};
18use spargebra::algebra::{GraphPattern, GraphTarget, QueryDataset};
19use spargebra::term::{
20 GraphNamePattern, GroundQuad, GroundQuadPattern, GroundTerm, GroundTermPattern,
21 NamedNodePattern, QuadPattern, TermPattern,
22};
23use spargebra::{GraphUpdateOperation, Update};
24use std::collections::{BTreeSet, HashMap};
25
26#[derive(Debug)]
28pub enum PlannedOp {
29 Sql(Vec<Statement>),
31 Fallback(usize, String),
33}
34
35fn bnode_map(t: &Term, map: &mut HashMap<BlankNode, BlankNode>) -> Term {
36 match t {
37 Term::BlankNode(b) => map.entry(b.clone()).or_default().clone().into(),
38 Term::Triple(tr) => oxrdf::Triple::new(
39 match &tr.subject {
40 NamedOrBlankNode::BlankNode(b) => {
41 NamedOrBlankNode::from(map.entry(b.clone()).or_default().clone())
42 }
43 s => s.clone(),
44 },
45 tr.predicate.clone(),
46 bnode_map(&tr.object, map),
47 )
48 .into(),
49 t => t.clone(),
50 }
51}
52
53pub fn insert_data_quads(data: &[spargebra::term::Quad]) -> Vec<Quad> {
55 let mut map: HashMap<BlankNode, BlankNode> = HashMap::new();
56 data.iter()
57 .map(|q| {
58 let s = match &q.subject {
59 NamedOrBlankNode::BlankNode(b) => {
60 NamedOrBlankNode::from(map.entry(b.clone()).or_default().clone())
61 }
62 s => s.clone(),
63 };
64 let g = match &q.graph_name {
65 spargebra::term::GraphName::NamedNode(n) => GraphName::NamedNode(n.clone()),
66 spargebra::term::GraphName::DefaultGraph => GraphName::DefaultGraph,
67 };
68 Quad::new(s, q.predicate.clone(), bnode_map(&q.object, &mut map), g)
69 })
70 .collect()
71}
72
73fn ground_term(t: &GroundTerm) -> Term {
74 crate::compiler::ground_to_term(t)
75}
76
77pub fn delete_data_quads(data: &[GroundQuad]) -> Vec<Quad> {
79 data.iter()
80 .map(|q| {
81 Quad::new(
82 q.subject.clone(),
83 q.predicate.clone(),
84 ground_term(&q.object),
85 match &q.graph_name {
86 spargebra::term::GraphName::NamedNode(n) => GraphName::NamedNode(n.clone()),
87 spargebra::term::GraphName::DefaultGraph => GraphName::DefaultGraph,
88 },
89 )
90 })
91 .collect()
92}
93
94fn guard(column: &str, condition: &str) -> Statement {
95 Statement::new(format!(
96 "INSERT INTO oxilite_guard({column}) SELECT 1 WHERE {condition}"
97 ))
98}
99
100fn clear_statements(
101 graph: &GraphTarget,
102 silent: bool,
103 drop: bool,
104 caps: &Capabilities,
105) -> Vec<Statement> {
106 let _ = caps;
107 let mut s = Vec::new();
108 match graph {
109 GraphTarget::NamedNode(g) => {
110 let id = named_node_id(g.as_str());
111 if !silent {
112 s.push(guard(
113 "graph_does_not_exist",
114 &format!("NOT EXISTS (SELECT 1 FROM graphs WHERE id = {id})"),
115 ));
116 }
117 s.push(Statement::new(format!("DELETE FROM quads WHERE g = {id}")));
118 if drop {
119 s.push(Statement::new(format!(
120 "DELETE FROM graphs WHERE id = {id}"
121 )));
122 }
123 }
124 GraphTarget::DefaultGraph => {
125 s.push(Statement::new(format!(
126 "DELETE FROM quads WHERE g = {DEFAULT_GRAPH_ID}"
127 )));
128 }
129 GraphTarget::NamedGraphs => {
130 s.push(Statement::new(format!(
131 "DELETE FROM quads WHERE g <> {DEFAULT_GRAPH_ID}"
132 )));
133 if drop {
134 s.push("DELETE FROM graphs".into());
135 }
136 }
137 GraphTarget::AllGraphs => {
138 s.push("DELETE FROM quads".into());
139 if drop {
140 s.push("DELETE FROM graphs".into());
141 }
142 }
143 }
144 s
145}
146
147enum Slot {
149 Sql(String),
150}
151
152#[allow(clippy::too_many_arguments)]
159fn delete_insert(
160 delete: &[GroundQuadPattern],
161 insert: &[QuadPattern],
162 using: Option<&QueryDataset>,
163 pattern: &GraphPattern,
164 base_iri: Option<String>,
165 stats: &Stats,
166 caps: &Capabilities,
167 options: &QueryOptions,
168) -> Result<Vec<Statement>> {
169 let mut c = Compiler::new(stats, caps, options, using, base_iri);
170 let block = c.pattern(pattern)?;
171 let mut vars: BTreeSet<Variable> = BTreeSet::new();
173 let mut bnodes: BTreeSet<String> = BTreeSet::new();
174 fn tvars(
175 t: &TermPattern,
176 vars: &mut BTreeSet<Variable>,
177 bnodes: &mut BTreeSet<String>,
178 ) -> Result<()> {
179 match t {
180 TermPattern::Variable(v) => {
181 vars.insert(v.clone());
182 }
183 TermPattern::BlankNode(b) => {
184 bnodes.insert(b.as_str().to_string());
185 }
186 TermPattern::Triple(_) => {
187 return Err(Error::unsupported("triple terms in update templates"))
188 }
189 _ => {}
190 }
191 Ok(())
192 }
193 fn gvars(t: &GroundTermPattern, vars: &mut BTreeSet<Variable>) -> Result<()> {
194 match t {
195 GroundTermPattern::Variable(v) => {
196 vars.insert(v.clone());
197 }
198 GroundTermPattern::Triple(_) => {
199 return Err(Error::unsupported("triple terms in update templates"))
200 }
201 _ => {}
202 }
203 Ok(())
204 }
205 let add_nn = |p: &NamedNodePattern, vars: &mut BTreeSet<Variable>| {
206 if let NamedNodePattern::Variable(v) = p {
207 vars.insert(v.clone());
208 }
209 };
210 let add_g = |g: &GraphNamePattern, vars: &mut BTreeSet<Variable>| {
211 if let GraphNamePattern::Variable(v) = g {
212 vars.insert(v.clone());
213 }
214 };
215 for q in delete {
216 gvars(&q.subject, &mut vars)?;
217 add_nn(&q.predicate, &mut vars);
218 gvars(&q.object, &mut vars)?;
219 add_g(&q.graph_name, &mut vars);
220 }
221 for q in insert {
222 tvars(&q.subject, &mut vars, &mut bnodes)?;
223 add_nn(&q.predicate, &mut vars);
224 tvars(&q.object, &mut vars, &mut bnodes)?;
225 add_g(&q.graph_name, &mut vars);
226 }
227 let idxs: Vec<usize> = vars.iter().map(|v| c.var(v)).collect();
228 let mut cols = HashMap::new();
229 let mut unstorable = Vec::new();
231 for (v, idx) in vars.iter().zip(&idxs) {
232 let key = match block.cols.get(idx).map(|b| &b.col) {
233 None => "NULL".to_string(),
234 Some(Col::Id(_)) => format!("w.v{idx}"),
235 Some(Col::Val(val)) if val.id.is_some() => format!("w.v{idx}_i"),
236 Some(Col::Val(val)) if val.stat == crate::compiler::expr::Stat::Numeric => {
237 let (n, t) = (format!("w.v{idx}_n"), format!("w.v{idx}_t"));
240 let slot = format!(
241 "(CASE WHEN {t} = 1 AND {n} = CAST({n} AS INTEGER) AND abs({n}) < 288230376151711744 THEN {} + CAST({n} AS INTEGER) END)",
242 crate::compiler::expr::INT_BASE
243 );
244 unstorable.push(format!("({n} IS NOT NULL AND {slot} IS NULL)"));
245 slot
246 }
247 Some(Col::Val(_)) => {
248 return Err(Error::unsupported(format!(
249 "computed value ?{} in an update template",
250 v.as_str()
251 )))
252 }
253 };
254 cols.insert(v.clone(), key);
255 }
256 let mut rows = EncodedRows::default();
257 let bnode_base = Tag::BlankNode.base();
258 let bnode_salt: HashMap<String, i64> = bnodes
259 .iter()
260 .map(|l| (l.clone(), crate::encoding::blank_node_id(l) & PAYLOAD_MASK))
261 .collect();
262 let term = |t: &Term, rows: &mut EncodedRows| rows.term(t.as_ref()).to_string();
263 let tslot = |t: &TermPattern, rows: &mut EncodedRows| -> Slot {
264 Slot::Sql(match t {
265 TermPattern::NamedNode(n) => term(&n.clone().into(), rows),
266 TermPattern::Literal(l) => term(&l.clone().into(), rows),
267 TermPattern::Variable(v) => cols[v].clone(),
268 TermPattern::BlankNode(b) => format!(
270 "({bnode_base} + ((w.rnd + {}) & {PAYLOAD_MASK}))",
271 bnode_salt[b.as_str()]
272 ),
273 TermPattern::Triple(_) => unreachable!("rejected above"),
274 })
275 };
276 let gslot = |t: &GroundTermPattern, rows: &mut EncodedRows| -> Slot {
277 Slot::Sql(match t {
278 GroundTermPattern::NamedNode(n) => term(&n.clone().into(), rows),
279 GroundTermPattern::Literal(l) => term(&l.clone().into(), rows),
280 GroundTermPattern::Variable(v) => cols[v].clone(),
281 GroundTermPattern::Triple(_) => unreachable!("rejected above"),
282 })
283 };
284 let nslot = |p: &NamedNodePattern, rows: &mut EncodedRows| -> String {
285 match p {
286 NamedNodePattern::NamedNode(n) => rows.iri(n.as_str()).to_string(),
287 NamedNodePattern::Variable(v) => cols[v].clone(),
288 }
289 };
290 let grslot = |g: &GraphNamePattern, rows: &mut EncodedRows| -> String {
291 match g {
292 GraphNamePattern::DefaultGraph => DEFAULT_GRAPH_ID.to_string(),
293 GraphNamePattern::NamedNode(n) => rows.iri(n.as_str()).to_string(),
294 GraphNamePattern::Variable(v) => cols[v].clone(),
295 }
296 };
297 let mut selects = Vec::new();
298 let valid_subject = |x: &str| format!("({x} >> {PAYLOAD_BITS}) IN (1, 2)");
299 let valid_iri = |x: &str| format!("({x} >> {PAYLOAD_BITS}) = 1");
300 let valid_graph = |x: &str| format!("({x} = 0 OR ({x} >> {PAYLOAD_BITS}) IN (1, 2))");
301 for q in delete {
302 let (Slot::Sql(s), p, Slot::Sql(o), g) = (
303 gslot(&q.subject, &mut rows),
304 nslot(&q.predicate, &mut rows),
305 gslot(&q.object, &mut rows),
306 grslot(&q.graph_name, &mut rows),
307 );
308 selects.push(format!(
309 "SELECT 0, {s}, {p}, {o}, {g} FROM w WHERE {s} IS NOT NULL AND {p} IS NOT NULL AND {o} IS NOT NULL AND {g} IS NOT NULL"
310 ));
311 }
312 for q in insert {
313 let (Slot::Sql(s), p, Slot::Sql(o), g) = (
314 tslot(&q.subject, &mut rows),
315 nslot(&q.predicate, &mut rows),
316 tslot(&q.object, &mut rows),
317 grslot(&q.graph_name, &mut rows),
318 );
319 selects.push(format!(
320 "SELECT 1, {s}, {p}, {o}, {g} FROM w WHERE {s} IS NOT NULL AND {p} IS NOT NULL AND {o} IS NOT NULL AND {g} IS NOT NULL AND {} AND {} AND {}",
321 valid_subject(&s),
322 valid_iri(&p),
323 valid_graph(&g)
324 ));
325 }
326 for t in c.constants.values() {
328 rows.term(t.as_ref());
329 }
330 rows.dedup();
331 let mut out = vec![Statement::new("DELETE FROM update_buffer")];
332 if !selects.is_empty() {
333 let mut where_block = block;
334 where_block
335 .extra_select
336 .push("(abs(random()) & 576460752303423487) AS rnd".into());
337 let where_sql = where_block.to_select(Some(&idxs), false);
338 if !unstorable.is_empty() {
339 out.push(Statement::new(format!(
340 "WITH w AS ({where_sql}) INSERT INTO oxilite_guard(computed_value_not_storable) SELECT 1 FROM w WHERE {} LIMIT 1",
341 unstorable.join(" OR ")
342 )));
343 }
344 out.push(Statement::new(format!(
345 "WITH w AS MATERIALIZED ({where_sql}) INSERT INTO update_buffer(op, s, p, o, g) {}",
346 crate::sql::union_all(selects, caps.max_compound_select)
347 )));
348 }
349 out.extend(term_statements(&rows, caps));
350 if !bnodes.is_empty() {
351 for c in ["s", "o"] {
353 out.push(Statement::new(format!(
354 "INSERT OR IGNORE INTO terms(id, lex) SELECT DISTINCT {c}, 'ox' || lower(hex({c})) FROM update_buffer WHERE op = 1 AND ({c} >> {PAYLOAD_BITS}) = 2 AND NOT EXISTS (SELECT 1 FROM terms t WHERE t.id = update_buffer.{c})"
355 )));
356 }
357 }
358 out.push(Statement::new(
359 "DELETE FROM quads WHERE (s, p, o, g) IN (SELECT s, p, o, g FROM update_buffer WHERE op = 0)",
360 ));
361 out.push(Statement::new(
362 if caps.versioning >= crate::version::Versioning::Stamped {
363 format!(
364 "INSERT OR IGNORE INTO quads(s, p, o, g, t) SELECT s, p, o, g, {} FROM update_buffer WHERE op = 1",
365 crate::version::CURRENT_TICK
366 )
367 } else {
368 "INSERT OR IGNORE INTO quads(s, p, o, g) SELECT s, p, o, g FROM update_buffer WHERE op = 1".to_owned()
369 },
370 ));
371 out.push(Statement::new(
372 "INSERT OR IGNORE INTO graphs(id) SELECT DISTINCT g FROM update_buffer WHERE op = 1 AND g <> 0",
373 ));
374 out.push(Statement::new("DELETE FROM update_buffer"));
375 Ok(out)
376}
377
378pub fn plan_update(update: &Update, caps: &Capabilities) -> Result<Vec<PlannedOp>> {
380 plan_update_with(update, &Stats::default(), caps, &QueryOptions::default())
381}
382
383pub fn plan_update_with(
385 update: &Update,
386 stats: &Stats,
387 caps: &Capabilities,
388 options: &QueryOptions,
389) -> Result<Vec<PlannedOp>> {
390 let mut out = Vec::new();
391 for (i, op) in update.operations.iter().enumerate() {
392 out.push(match op {
393 GraphUpdateOperation::InsertData { data } => {
394 let quads = insert_data_quads(data);
395 PlannedOp::Sql(
396 EncodedQuads::new(quads.iter().map(Quad::as_ref)).insert_statements(caps),
397 )
398 }
399 GraphUpdateOperation::DeleteData { data } => {
400 let quads = delete_data_quads(data);
401 if quads.is_empty() {
402 PlannedOp::Sql(Vec::new())
403 } else {
404 PlannedOp::Sql(
405 EncodedQuads::new(quads.iter().map(Quad::as_ref)).delete_statements(caps),
406 )
407 }
408 }
409 GraphUpdateOperation::Clear { silent, graph } => {
410 PlannedOp::Sql(clear_statements(graph, *silent, false, caps))
411 }
412 GraphUpdateOperation::Drop { silent, graph } => {
413 PlannedOp::Sql(clear_statements(graph, *silent, true, caps))
414 }
415 GraphUpdateOperation::Create { silent, graph } => {
416 let mut rows = EncodedRows::default();
417 let id = rows.iri(graph.as_str());
418 let mut s = term_statements(&rows, caps);
419 if *silent {
420 s.push(Statement::new(format!(
421 "INSERT OR IGNORE INTO graphs(id) VALUES ({id})"
422 )));
423 } else {
424 s.push(guard(
425 "graph_already_exists",
426 &format!("EXISTS (SELECT 1 FROM graphs WHERE id = {id})"),
427 ));
428 s.push(Statement::new(format!(
429 "INSERT INTO graphs(id) VALUES ({id})"
430 )));
431 }
432 PlannedOp::Sql(s)
433 }
434 GraphUpdateOperation::Load { silent, .. } => {
435 if *silent {
436 PlannedOp::Sql(Vec::new())
437 } else {
438 return Err(Error::unsupported("LOAD (the core has no network access)"));
439 }
440 }
441 GraphUpdateOperation::DeleteInsert {
442 delete,
443 insert,
444 using,
445 pattern,
446 } => match delete_insert(
447 delete,
448 insert,
449 using.as_ref(),
450 pattern,
451 update.base_iri.as_ref().map(|b| b.as_str().to_string()),
452 stats,
453 caps,
454 options,
455 ) {
456 Ok(stmts) => PlannedOp::Sql(stmts),
457 Err(e) if e.is_unsupported() => {
458 PlannedOp::Fallback(i, format!("DELETE/INSERT … WHERE ({e})"))
459 }
460 Err(e) => return Err(e),
461 },
462 });
463 }
464 if crate::reason::update_touches_schema(update) {
466 out.push(PlannedOp::Sql(crate::reason::closure_statements()));
467 }
468 if crate::shapes::update_touches_shapes(update) {
469 out.push(PlannedOp::Sql(crate::shapes::refresh_statements()));
470 }
471 Ok(out)
472}
473
474pub fn explain_plan(plan: &[PlannedOp]) -> String {
476 let mut out = String::new();
477 for (i, p) in plan.iter().enumerate() {
478 match p {
479 PlannedOp::Sql(s) => {
480 out.push_str(&format!(
481 "-- operation {i}: compiled to {} SQL statement(s)\n",
482 s.len()
483 ));
484 for st in s {
485 out.push_str(&st.sql);
486 out.push_str(";\n");
487 }
488 }
489 PlannedOp::Fallback(_, why) => {
490 out.push_str(&format!("-- operation {i}: not compiled ({why}); evaluated by the fallback (sync backends only)\n"));
491 }
492 }
493 }
494 out
495}
496
497pub fn graph_name_id(g: &GraphName) -> i64 {
499 graph_id(g.as_ref())
500}