Skip to main content

oxilite_core/
fallback.rs

1//! Fallback evaluation with `spareval` for queries the SQL compiler cannot express.
2//!
3//! Only available on synchronous backends: `spareval` pulls quads through an iterator
4//! interface. Hash ids make `internalize_term` free; `externalize_term` is a cached lookup.
5//!
6// @lat: [[architecture#SPARQL to SQL compiler#Fallback evaluator]]
7
8use crate::encoding::{
9    decode_inline, decode_row, make_triple, tag_of, term_id, Tag, DEFAULT_GRAPH_ID,
10};
11use crate::error::{Error, Result};
12use crate::job::SyncBackend;
13use crate::query::QueryOutput;
14use crate::sql::{Request, SqlValue, Statement};
15use oxrdf::Term;
16use spareval::{InternalQuad, QueryEvaluator, QueryResults, QueryableDataset};
17use spargebra::Query;
18use std::cell::RefCell;
19use std::collections::{BTreeSet, HashMap};
20
21/// A `spareval` dataset over a SQL backend.
22pub struct SqlDataset<'a, B: SyncBackend> {
23    backend: &'a B,
24    cache: RefCell<HashMap<i64, Term>>,
25    /// Query-time reasoning (quads are read from the entailed-triple source).
26    reasoning: crate::reason::Reasoning,
27    inferred: bool,
28    hide_schema: bool,
29    transitive: BTreeSet<i64>,
30    scopes: BTreeSet<i64>,
31    /// Read the store as it was at this tick (see `version`).
32    as_of: Option<i64>,
33}
34
35impl<'a, B: SyncBackend> SqlDataset<'a, B> {
36    pub fn new(backend: &'a B) -> Self {
37        Self {
38            backend,
39            cache: RefCell::new(HashMap::new()),
40            reasoning: crate::reason::Reasoning::None,
41            inferred: false,
42            hide_schema: false,
43            transitive: BTreeSet::new(),
44            scopes: BTreeSet::new(),
45            as_of: None,
46        }
47    }
48
49    /// Reads entailed triples, as the SQL compiler does for these options.
50    pub fn with_options(mut self, options: &crate::QueryOptions) -> Result<Self> {
51        self.reasoning = options.reasoning;
52        self.inferred = options.include_inferred;
53        self.hide_schema = !options.include_schema_graphs;
54        self.as_of = options.as_of_tick;
55        if self.reasoning == crate::reason::Reasoning::OwlQl {
56            let stmt = crate::reason::transitive_statement(|c| self.id_col(c));
57            for row in self.rows(stmt.sql)? {
58                if let Some(p) = row.first().and_then(SqlValue::as_i64) {
59                    self.transitive.insert(p);
60                }
61            }
62        }
63        if self.reasoning != crate::reason::Reasoning::None {
64            let stmt = crate::registry::scopes_statement(|c| self.id_col(c));
65            for row in self.rows(stmt.sql)? {
66                if let Some(g) = row.first().and_then(SqlValue::as_i64) {
67                    self.scopes.insert(g);
68                }
69            }
70        }
71        Ok(self)
72    }
73
74    fn id_col(&self, c: &str) -> String {
75        if self.backend.capabilities().int64_as_text {
76            format!("CAST({c} AS TEXT)")
77        } else {
78            c.into()
79        }
80    }
81
82    fn rows(&self, sql: String) -> Result<Vec<Vec<SqlValue>>> {
83        let mut r = self
84            .backend
85            .execute(&Request::read(vec![Statement::new(sql)]))?;
86        Ok(r.pop().map(|rs| rs.rows).unwrap_or_default())
87    }
88
89    pub fn lookup(&self, id: i64) -> Result<Term> {
90        if let Some(t) = self.cache.borrow().get(&id) {
91            return Ok(t.clone());
92        }
93        if let Some(t) = decode_inline(id) {
94            return Ok(t);
95        }
96        let t = if tag_of(id) == Some(Tag::Triple) {
97            let rows = self.rows(format!(
98                "SELECT {}, {}, {} FROM triple_terms WHERE id = {id}",
99                self.id_col("s"),
100                self.id_col("p"),
101                self.id_col("o")
102            ))?;
103            let row = rows
104                .first()
105                .ok_or_else(|| Error::corrupted(format!("triple term {id} missing")))?;
106            let get = |i: usize| {
107                row.get(i)
108                    .and_then(SqlValue::as_i64)
109                    .ok_or_else(|| Error::corrupted("bad triple term"))
110            };
111            Term::from(make_triple(
112                self.lookup(get(0)?)?,
113                self.lookup(get(1)?)?,
114                self.lookup(get(2)?)?,
115            )?)
116        } else {
117            let rows = self.rows(format!(
118                "SELECT lex, dt, lang, dir FROM terms WHERE id = {id}"
119            ))?;
120            let mut row = rows
121                .into_iter()
122                .next()
123                .ok_or_else(|| Error::corrupted(format!("term {id} missing")))?
124                .into_iter();
125            let lex = row
126                .next()
127                .and_then(SqlValue::into_string)
128                .unwrap_or_default();
129            let dt = row.next().and_then(SqlValue::into_string);
130            let lang = row.next().and_then(SqlValue::into_string);
131            let dir = row.next().and_then(|v| v.as_i64());
132            decode_row(id, lex, dt, lang, dir)?
133        };
134        self.cache.borrow_mut().insert(id, t.clone());
135        Ok(t)
136    }
137}
138
139impl<'a, B: SyncBackend> QueryableDataset<'a> for SqlDataset<'a, B> {
140    type InternalTerm = i64;
141    type Error = Error;
142
143    fn internal_quads_for_pattern(
144        &self,
145        subject: Option<&i64>,
146        predicate: Option<&i64>,
147        object: Option<&i64>,
148        graph_name: Option<Option<&i64>>,
149    ) -> impl Iterator<Item = Result<InternalQuad<i64>>> + use<'a, B> {
150        let mut w = Vec::new();
151        if let Some(s) = subject {
152            w.push(format!("s = {s}"));
153        }
154        if let Some(p) = predicate {
155            w.push(format!("p = {p}"));
156        }
157        if let Some(o) = object {
158            w.push(format!("o = {o}"));
159        }
160        match graph_name {
161            None => w.push(format!("g <> {DEFAULT_GRAPH_ID}")),
162            Some(None) => w.push(format!("g = {DEFAULT_GRAPH_ID}")),
163            Some(Some(g)) => w.push(format!("g = {g}")),
164        }
165        let ent = crate::reason::Entailment {
166            reasoning: self.reasoning,
167            inferred: self.inferred,
168            hide_schema: self.hide_schema,
169            as_of: self.as_of,
170            transitive: &self.transitive,
171            scopes: &self.scopes,
172            max_compound: self.backend.capabilities().max_compound_select,
173        };
174        let source = if ent.active() {
175            format!(
176                "{} AS quads",
177                ent.source(
178                    subject.copied(),
179                    predicate.copied(),
180                    object.copied(),
181                    &crate::reason::GraphFilter::Keep
182                )
183            )
184        } else {
185            format!("{} AS quads", ent.base())
186        };
187        let sql = format!(
188            "SELECT {}, {}, {}, {} FROM {source} WHERE {}",
189            self.id_col("s"),
190            self.id_col("p"),
191            self.id_col("o"),
192            self.id_col("g"),
193            w.join(" AND ")
194        );
195        let result: Vec<Result<InternalQuad<i64>>> = match self.rows(sql) {
196            Ok(rows) => rows
197                .into_iter()
198                .map(|row| {
199                    let ids: Vec<i64> = row.iter().filter_map(SqlValue::as_i64).collect();
200                    let [s, p, o, g] = ids[..] else {
201                        return Err(Error::corrupted("bad quad row"));
202                    };
203                    Ok(InternalQuad {
204                        subject: s,
205                        predicate: p,
206                        object: o,
207                        graph_name: (g != DEFAULT_GRAPH_ID).then_some(g),
208                    })
209                })
210                .collect(),
211            Err(e) => vec![Err(e)],
212        };
213        result.into_iter()
214    }
215
216    fn internal_named_graphs(&self) -> impl Iterator<Item = Result<i64>> + use<'a, B> {
217        let r: Vec<Result<i64>> =
218            match self.rows(format!("SELECT {} FROM graphs", self.id_col("id"))) {
219                Ok(rows) => rows
220                    .into_iter()
221                    .filter_map(|r| r.first().and_then(SqlValue::as_i64))
222                    .map(Ok)
223                    .collect(),
224                Err(e) => vec![Err(e)],
225            };
226        r.into_iter()
227    }
228
229    fn contains_internal_graph_name(&self, graph_name: &i64) -> Result<bool> {
230        let rows = self.rows(format!(
231            "SELECT EXISTS (SELECT 1 FROM graphs WHERE id = {graph_name}) OR EXISTS (SELECT 1 FROM quads WHERE g = {graph_name})"
232        ))?;
233        Ok(rows
234            .first()
235            .and_then(|r| r.first())
236            .and_then(SqlValue::as_i64)
237            == Some(1))
238    }
239
240    fn internalize_term(&self, term: Term) -> Result<i64> {
241        let id = term_id(term.as_ref());
242        self.cache.borrow_mut().entry(id).or_insert(term);
243        Ok(id)
244    }
245
246    fn externalize_term(&self, term: i64) -> Result<Term> {
247        self.lookup(term)
248    }
249}
250
251/// Applies the dataset options (union default graph, explicit default / named graphs) to a
252/// spareval dataset specification.
253pub fn apply_dataset_options(
254    spec: &mut spareval::QueryDatasetSpecification,
255    query: &Query,
256    options: &crate::QueryOptions,
257    lookup: impl Fn(i64) -> Option<oxrdf::Term>,
258) {
259    if options.union_default_graph && query_dataset(query).is_none() {
260        spec.set_default_graph_as_union();
261    }
262    let graph = |id: &i64| -> Option<oxrdf::GraphName> {
263        if *id == DEFAULT_GRAPH_ID {
264            return Some(oxrdf::GraphName::DefaultGraph);
265        }
266        match lookup(*id)? {
267            Term::NamedNode(n) => Some(n.into()),
268            Term::BlankNode(b) => Some(b.into()),
269            _ => None,
270        }
271    };
272    if let Some(d) = &options.default_graph {
273        spec.set_default_graph(d.iter().filter_map(graph).collect());
274    }
275    if let Some(n) = &options.named_graphs {
276        spec.set_available_named_graphs(
277            n.iter()
278                .filter_map(|id| match graph(id)? {
279                    oxrdf::GraphName::NamedNode(n) => Some(n.into()),
280                    oxrdf::GraphName::BlankNode(b) => Some(b.into()),
281                    oxrdf::GraphName::DefaultGraph => None,
282                })
283                .collect(),
284        );
285    }
286}
287
288/// Evaluates a query with `spareval` over a sync backend.
289pub fn evaluate<B: SyncBackend>(
290    backend: &B,
291    query: &Query,
292    options: &crate::QueryOptions,
293) -> Result<QueryOutput> {
294    let evaluator = options.functions.install(
295        QueryEvaluator::new()
296            .with_custom_function(crate::text::text_match_name(), crate::text::text_match),
297    );
298    let mut prepared = evaluator.prepare(query);
299    let dataset = SqlDataset::new(backend).with_options(options)?;
300    apply_dataset_options(prepared.dataset_mut(), query, options, |id| {
301        dataset.lookup(id).ok()
302    });
303    Ok(match prepared.execute(dataset)? {
304        QueryResults::Solutions(solutions) => {
305            let variables = solutions.variables().to_vec();
306            let mut rows = Vec::new();
307            for s in solutions {
308                let s = s?;
309                rows.push(variables.iter().map(|v| s.get(v).cloned()).collect());
310            }
311            QueryOutput::Solutions { variables, rows }
312        }
313        QueryResults::Boolean(b) => QueryOutput::Boolean(b),
314        QueryResults::Graph(triples) => QueryOutput::Graph(triples.collect::<Result<Vec<_>, _>>()?),
315    })
316}
317
318fn query_dataset(q: &Query) -> Option<&spargebra::algebra::QueryDataset> {
319    match q {
320        Query::Select { dataset, .. }
321        | Query::Construct { dataset, .. }
322        | Query::Describe { dataset, .. }
323        | Query::Ask { dataset, .. } => dataset.as_ref(),
324    }
325}
326
327/// Evaluates a `DELETE/INSERT … WHERE` operation with `spareval`, returning the quads to
328/// delete and to insert (both computed before any modification).
329pub fn delete_insert<B: SyncBackend>(
330    backend: &B,
331    op: &spargebra::GraphUpdateOperation,
332    base_iri: Option<&oxiri::Iri<String>>,
333) -> Result<(Vec<oxrdf::Quad>, Vec<oxrdf::Quad>)> {
334    delete_insert_with(
335        backend,
336        op,
337        base_iri,
338        &crate::functions::Functions::default(),
339    )
340}
341
342/// [`delete_insert`] with host functions callable from the `WHERE` clause.
343pub fn delete_insert_with<B: SyncBackend>(
344    backend: &B,
345    op: &spargebra::GraphUpdateOperation,
346    base_iri: Option<&oxiri::Iri<String>>,
347    functions: &crate::functions::Functions,
348) -> Result<(Vec<oxrdf::Quad>, Vec<oxrdf::Quad>)> {
349    let spargebra::GraphUpdateOperation::DeleteInsert {
350        delete,
351        insert,
352        using,
353        pattern,
354    } = op
355    else {
356        return Err(Error::Other("not a DELETE/INSERT operation".into()));
357    };
358    let evaluator = functions.install(
359        QueryEvaluator::new()
360            .with_custom_function(crate::text::text_match_name(), crate::text::text_match),
361    );
362    let prepared = evaluator.prepare_delete_insert(
363        delete.clone(),
364        insert.clone(),
365        base_iri.cloned(),
366        using.clone(),
367        pattern,
368    );
369    let mut deletes = Vec::new();
370    let mut inserts = Vec::new();
371    for q in prepared.execute(SqlDataset::new(backend))? {
372        match q? {
373            spareval::DeleteInsertQuad::Delete(q) => deletes.push(q),
374            spareval::DeleteInsertQuad::Insert(q) => inserts.push(q),
375        }
376    }
377    Ok((deletes, inserts))
378}