1use 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
21pub struct SqlDataset<'a, B: SyncBackend> {
23 backend: &'a B,
24 cache: RefCell<HashMap<i64, Term>>,
25 reasoning: crate::reason::Reasoning,
27 inferred: bool,
28 hide_schema: bool,
29 transitive: BTreeSet<i64>,
30 scopes: BTreeSet<i64>,
31 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 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
251pub 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
288pub 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
327pub 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
342pub 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}