Skip to main content

oxilite_core/
update.rs

1//! SPARQL UPDATE planning.
2//!
3//! Operations become SQL statements appended to one atomic request, so a whole update
4//! request is applied as a single transaction (one D1 batch). Operations that cannot be
5//! compiled are reported as [`PlannedOp::Fallback`] for sync backends.
6//!
7// @lat: [[architecture#Updates and atomicity]]
8
9use 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/// A planned update operation.
27#[derive(Debug)]
28pub enum PlannedOp {
29    /// SQL statements, to run in order inside the update's transaction.
30    Sql(Vec<Statement>),
31    /// Operation `index` of the update must be evaluated by the fallback (sync backends).
32    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
53/// `INSERT DATA` quads with fresh blank nodes.
54pub 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
77/// Converts `DELETE DATA` quads.
78pub 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
147/// A template position after compilation: a SQL expression of the term id.
148enum Slot {
149    Sql(String),
150}
151
152/// Compiles `DELETE { … } INSERT { … } WHERE { … }` into self-reading statements: the WHERE
153/// solutions are evaluated once (a MATERIALIZED CTE) into `update_buffer`, then deletions and
154/// insertions are applied from the buffer — SPARQL's "evaluate, then apply" semantics inside
155/// one atomic batch.
156///
157// @lat: [[architecture#Updates and atomicity]]
158#[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    // Variables used by the templates.
172    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    // Rows whose computed value cannot be stored from SQL (not an inline integer).
230    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                // Computed numbers are stored as inline integer ids; anything else aborts the
238                // batch (native stores then use the fallback).
239                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            // A fresh blank node per solution: the row's random seed mixed with the label.
269            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    // Constants that may end up stored need their dictionary rows.
327    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        // Dictionary rows for the fresh blank nodes.
352        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        "INSERT OR IGNORE INTO quads(s, p, o, g) SELECT s, p, o, g FROM update_buffer WHERE op = 1",
363    ));
364    out.push(Statement::new(
365        "INSERT OR IGNORE INTO graphs(id) SELECT DISTINCT g FROM update_buffer WHERE op = 1 AND g <> 0",
366    ));
367    out.push(Statement::new("DELETE FROM update_buffer"));
368    Ok(out)
369}
370
371/// Plans a SPARQL update.
372pub fn plan_update(update: &Update, caps: &Capabilities) -> Result<Vec<PlannedOp>> {
373    plan_update_with(update, &Stats::default(), caps, &QueryOptions::default())
374}
375
376/// Plans a SPARQL update using planner statistics.
377pub fn plan_update_with(
378    update: &Update,
379    stats: &Stats,
380    caps: &Capabilities,
381    options: &QueryOptions,
382) -> Result<Vec<PlannedOp>> {
383    let mut out = Vec::new();
384    for (i, op) in update.operations.iter().enumerate() {
385        out.push(match op {
386            GraphUpdateOperation::InsertData { data } => {
387                let quads = insert_data_quads(data);
388                PlannedOp::Sql(
389                    EncodedQuads::new(quads.iter().map(Quad::as_ref)).insert_statements(caps),
390                )
391            }
392            GraphUpdateOperation::DeleteData { data } => {
393                let quads = delete_data_quads(data);
394                if quads.is_empty() {
395                    PlannedOp::Sql(Vec::new())
396                } else {
397                    PlannedOp::Sql(
398                        EncodedQuads::new(quads.iter().map(Quad::as_ref)).delete_statements(caps),
399                    )
400                }
401            }
402            GraphUpdateOperation::Clear { silent, graph } => {
403                PlannedOp::Sql(clear_statements(graph, *silent, false, caps))
404            }
405            GraphUpdateOperation::Drop { silent, graph } => {
406                PlannedOp::Sql(clear_statements(graph, *silent, true, caps))
407            }
408            GraphUpdateOperation::Create { silent, graph } => {
409                let mut rows = EncodedRows::default();
410                let id = rows.iri(graph.as_str());
411                let mut s = term_statements(&rows, caps);
412                if *silent {
413                    s.push(Statement::new(format!(
414                        "INSERT OR IGNORE INTO graphs(id) VALUES ({id})"
415                    )));
416                } else {
417                    s.push(guard(
418                        "graph_already_exists",
419                        &format!("EXISTS (SELECT 1 FROM graphs WHERE id = {id})"),
420                    ));
421                    s.push(Statement::new(format!(
422                        "INSERT INTO graphs(id) VALUES ({id})"
423                    )));
424                }
425                PlannedOp::Sql(s)
426            }
427            GraphUpdateOperation::Load { silent, .. } => {
428                if *silent {
429                    PlannedOp::Sql(Vec::new())
430                } else {
431                    return Err(Error::unsupported("LOAD (the core has no network access)"));
432                }
433            }
434            GraphUpdateOperation::DeleteInsert {
435                delete,
436                insert,
437                using,
438                pattern,
439            } => match delete_insert(
440                delete,
441                insert,
442                using.as_ref(),
443                pattern,
444                update.base_iri.as_ref().map(|b| b.as_str().to_string()),
445                stats,
446                caps,
447                options,
448            ) {
449                Ok(stmts) => PlannedOp::Sql(stmts),
450                Err(e) if e.is_unsupported() => {
451                    PlannedOp::Fallback(i, format!("DELETE/INSERT … WHERE ({e})"))
452                }
453                Err(e) => return Err(e),
454            },
455        });
456    }
457    // Writes to schema triples keep the reasoning closure current, in the same transaction.
458    if crate::reason::update_touches_schema(update) {
459        out.push(PlannedOp::Sql(crate::reason::closure_statements()));
460    }
461    Ok(out)
462}
463
464/// Describes how each operation of an update runs (for `explain_update()`).
465pub fn explain_plan(plan: &[PlannedOp]) -> String {
466    let mut out = String::new();
467    for (i, p) in plan.iter().enumerate() {
468        match p {
469            PlannedOp::Sql(s) => {
470                out.push_str(&format!(
471                    "-- operation {i}: compiled to {} SQL statement(s)\n",
472                    s.len()
473                ));
474                for st in s {
475                    out.push_str(&st.sql);
476                    out.push_str(";\n");
477                }
478            }
479            PlannedOp::Fallback(_, why) => {
480                out.push_str(&format!("-- operation {i}: not compiled ({why}); evaluated by the fallback (sync backends only)\n"));
481            }
482        }
483    }
484    out
485}
486
487/// Graph id helper used by drivers.
488pub fn graph_name_id(g: &GraphName) -> i64 {
489    graph_id(g.as_ref())
490}