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 = QueryEvaluator::new()
295 .with_custom_function(crate::text::text_match_name(), crate::text::text_match);
296 let mut prepared = evaluator.prepare(query);
297 let dataset = SqlDataset::new(backend).with_options(options)?;
298 apply_dataset_options(prepared.dataset_mut(), query, options, |id| {
299 dataset.lookup(id).ok()
300 });
301 Ok(match prepared.execute(dataset)? {
302 QueryResults::Solutions(solutions) => {
303 let variables = solutions.variables().to_vec();
304 let mut rows = Vec::new();
305 for s in solutions {
306 let s = s?;
307 rows.push(variables.iter().map(|v| s.get(v).cloned()).collect());
308 }
309 QueryOutput::Solutions { variables, rows }
310 }
311 QueryResults::Boolean(b) => QueryOutput::Boolean(b),
312 QueryResults::Graph(triples) => QueryOutput::Graph(triples.collect::<Result<Vec<_>, _>>()?),
313 })
314}
315
316fn query_dataset(q: &Query) -> Option<&spargebra::algebra::QueryDataset> {
317 match q {
318 Query::Select { dataset, .. }
319 | Query::Construct { dataset, .. }
320 | Query::Describe { dataset, .. }
321 | Query::Ask { dataset, .. } => dataset.as_ref(),
322 }
323}
324
325pub fn delete_insert<B: SyncBackend>(
328 backend: &B,
329 op: &spargebra::GraphUpdateOperation,
330 base_iri: Option<&oxiri::Iri<String>>,
331) -> Result<(Vec<oxrdf::Quad>, Vec<oxrdf::Quad>)> {
332 let spargebra::GraphUpdateOperation::DeleteInsert {
333 delete,
334 insert,
335 using,
336 pattern,
337 } = op
338 else {
339 return Err(Error::Other("not a DELETE/INSERT operation".into()));
340 };
341 let evaluator = QueryEvaluator::new()
342 .with_custom_function(crate::text::text_match_name(), crate::text::text_match);
343 let prepared = evaluator.prepare_delete_insert(
344 delete.clone(),
345 insert.clone(),
346 base_iri.cloned(),
347 using.clone(),
348 pattern,
349 );
350 let mut deletes = Vec::new();
351 let mut inserts = Vec::new();
352 for q in prepared.execute(SqlDataset::new(backend))? {
353 match q? {
354 spareval::DeleteInsertQuad::Delete(q) => deletes.push(q),
355 spareval::DeleteInsertQuad::Insert(q) => inserts.push(q),
356 }
357 }
358 Ok((deletes, inserts))
359}