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 as_of: Option<i64>,
32}
33
34impl<'a, B: SyncBackend> SqlDataset<'a, B> {
35 pub fn new(backend: &'a B) -> Self {
36 Self {
37 backend,
38 cache: RefCell::new(HashMap::new()),
39 reasoning: crate::reason::Reasoning::None,
40 inferred: false,
41 hide_schema: false,
42 transitive: BTreeSet::new(),
43 as_of: None,
44 }
45 }
46
47 pub fn with_options(mut self, options: &crate::QueryOptions) -> Result<Self> {
49 self.reasoning = options.reasoning;
50 self.inferred = options.include_inferred;
51 self.hide_schema = !options.include_schema_graphs;
52 self.as_of = options.as_of_tick;
53 if self.reasoning == crate::reason::Reasoning::OwlQl {
54 let stmt = crate::reason::transitive_statement(|c| self.id_col(c));
55 for row in self.rows(stmt.sql)? {
56 if let Some(p) = row.first().and_then(SqlValue::as_i64) {
57 self.transitive.insert(p);
58 }
59 }
60 }
61 Ok(self)
62 }
63
64 fn id_col(&self, c: &str) -> String {
65 if self.backend.capabilities().int64_as_text {
66 format!("CAST({c} AS TEXT)")
67 } else {
68 c.into()
69 }
70 }
71
72 fn rows(&self, sql: String) -> Result<Vec<Vec<SqlValue>>> {
73 let mut r = self
74 .backend
75 .execute(&Request::read(vec![Statement::new(sql)]))?;
76 Ok(r.pop().map(|rs| rs.rows).unwrap_or_default())
77 }
78
79 pub fn lookup(&self, id: i64) -> Result<Term> {
80 if let Some(t) = self.cache.borrow().get(&id) {
81 return Ok(t.clone());
82 }
83 if let Some(t) = decode_inline(id) {
84 return Ok(t);
85 }
86 let t = if tag_of(id) == Some(Tag::Triple) {
87 let rows = self.rows(format!(
88 "SELECT {}, {}, {} FROM triple_terms WHERE id = {id}",
89 self.id_col("s"),
90 self.id_col("p"),
91 self.id_col("o")
92 ))?;
93 let row = rows
94 .first()
95 .ok_or_else(|| Error::corrupted(format!("triple term {id} missing")))?;
96 let get = |i: usize| {
97 row.get(i)
98 .and_then(SqlValue::as_i64)
99 .ok_or_else(|| Error::corrupted("bad triple term"))
100 };
101 Term::from(make_triple(
102 self.lookup(get(0)?)?,
103 self.lookup(get(1)?)?,
104 self.lookup(get(2)?)?,
105 )?)
106 } else {
107 let rows = self.rows(format!(
108 "SELECT lex, dt, lang, dir FROM terms WHERE id = {id}"
109 ))?;
110 let mut row = rows
111 .into_iter()
112 .next()
113 .ok_or_else(|| Error::corrupted(format!("term {id} missing")))?
114 .into_iter();
115 let lex = row
116 .next()
117 .and_then(SqlValue::into_string)
118 .unwrap_or_default();
119 let dt = row.next().and_then(SqlValue::into_string);
120 let lang = row.next().and_then(SqlValue::into_string);
121 let dir = row.next().and_then(|v| v.as_i64());
122 decode_row(id, lex, dt, lang, dir)?
123 };
124 self.cache.borrow_mut().insert(id, t.clone());
125 Ok(t)
126 }
127}
128
129impl<'a, B: SyncBackend> QueryableDataset<'a> for SqlDataset<'a, B> {
130 type InternalTerm = i64;
131 type Error = Error;
132
133 fn internal_quads_for_pattern(
134 &self,
135 subject: Option<&i64>,
136 predicate: Option<&i64>,
137 object: Option<&i64>,
138 graph_name: Option<Option<&i64>>,
139 ) -> impl Iterator<Item = Result<InternalQuad<i64>>> + use<'a, B> {
140 let mut w = Vec::new();
141 if let Some(s) = subject {
142 w.push(format!("s = {s}"));
143 }
144 if let Some(p) = predicate {
145 w.push(format!("p = {p}"));
146 }
147 if let Some(o) = object {
148 w.push(format!("o = {o}"));
149 }
150 match graph_name {
151 None => w.push(format!("g <> {DEFAULT_GRAPH_ID}")),
152 Some(None) => w.push(format!("g = {DEFAULT_GRAPH_ID}")),
153 Some(Some(g)) => w.push(format!("g = {g}")),
154 }
155 let ent = crate::reason::Entailment {
156 reasoning: self.reasoning,
157 inferred: self.inferred,
158 hide_schema: self.hide_schema,
159 as_of: self.as_of,
160 transitive: &self.transitive,
161 max_compound: self.backend.capabilities().max_compound_select,
162 };
163 let source = if ent.active() {
164 format!(
165 "{} AS quads",
166 ent.source(
167 subject.copied(),
168 predicate.copied(),
169 object.copied(),
170 &crate::reason::GraphFilter::Keep
171 )
172 )
173 } else {
174 format!("{} AS quads", ent.base())
175 };
176 let sql = format!(
177 "SELECT {}, {}, {}, {} FROM {source} WHERE {}",
178 self.id_col("s"),
179 self.id_col("p"),
180 self.id_col("o"),
181 self.id_col("g"),
182 w.join(" AND ")
183 );
184 let result: Vec<Result<InternalQuad<i64>>> = match self.rows(sql) {
185 Ok(rows) => rows
186 .into_iter()
187 .map(|row| {
188 let ids: Vec<i64> = row.iter().filter_map(SqlValue::as_i64).collect();
189 let [s, p, o, g] = ids[..] else {
190 return Err(Error::corrupted("bad quad row"));
191 };
192 Ok(InternalQuad {
193 subject: s,
194 predicate: p,
195 object: o,
196 graph_name: (g != DEFAULT_GRAPH_ID).then_some(g),
197 })
198 })
199 .collect(),
200 Err(e) => vec![Err(e)],
201 };
202 result.into_iter()
203 }
204
205 fn internal_named_graphs(&self) -> impl Iterator<Item = Result<i64>> + use<'a, B> {
206 let r: Vec<Result<i64>> =
207 match self.rows(format!("SELECT {} FROM graphs", self.id_col("id"))) {
208 Ok(rows) => rows
209 .into_iter()
210 .filter_map(|r| r.first().and_then(SqlValue::as_i64))
211 .map(Ok)
212 .collect(),
213 Err(e) => vec![Err(e)],
214 };
215 r.into_iter()
216 }
217
218 fn contains_internal_graph_name(&self, graph_name: &i64) -> Result<bool> {
219 let rows = self.rows(format!(
220 "SELECT EXISTS (SELECT 1 FROM graphs WHERE id = {graph_name}) OR EXISTS (SELECT 1 FROM quads WHERE g = {graph_name})"
221 ))?;
222 Ok(rows
223 .first()
224 .and_then(|r| r.first())
225 .and_then(SqlValue::as_i64)
226 == Some(1))
227 }
228
229 fn internalize_term(&self, term: Term) -> Result<i64> {
230 let id = term_id(term.as_ref());
231 self.cache.borrow_mut().entry(id).or_insert(term);
232 Ok(id)
233 }
234
235 fn externalize_term(&self, term: i64) -> Result<Term> {
236 self.lookup(term)
237 }
238}
239
240pub fn apply_dataset_options(
243 spec: &mut spareval::QueryDatasetSpecification,
244 query: &Query,
245 options: &crate::QueryOptions,
246 lookup: impl Fn(i64) -> Option<oxrdf::Term>,
247) {
248 if options.union_default_graph && query_dataset(query).is_none() {
249 spec.set_default_graph_as_union();
250 }
251 let graph = |id: &i64| -> Option<oxrdf::GraphName> {
252 if *id == DEFAULT_GRAPH_ID {
253 return Some(oxrdf::GraphName::DefaultGraph);
254 }
255 match lookup(*id)? {
256 Term::NamedNode(n) => Some(n.into()),
257 Term::BlankNode(b) => Some(b.into()),
258 _ => None,
259 }
260 };
261 if let Some(d) = &options.default_graph {
262 spec.set_default_graph(d.iter().filter_map(graph).collect());
263 }
264 if let Some(n) = &options.named_graphs {
265 spec.set_available_named_graphs(
266 n.iter()
267 .filter_map(|id| match graph(id)? {
268 oxrdf::GraphName::NamedNode(n) => Some(n.into()),
269 oxrdf::GraphName::BlankNode(b) => Some(b.into()),
270 oxrdf::GraphName::DefaultGraph => None,
271 })
272 .collect(),
273 );
274 }
275}
276
277pub fn evaluate<B: SyncBackend>(
279 backend: &B,
280 query: &Query,
281 options: &crate::QueryOptions,
282) -> Result<QueryOutput> {
283 let evaluator = QueryEvaluator::new()
284 .with_custom_function(crate::text::text_match_name(), crate::text::text_match);
285 let mut prepared = evaluator.prepare(query);
286 let dataset = SqlDataset::new(backend).with_options(options)?;
287 apply_dataset_options(prepared.dataset_mut(), query, options, |id| {
288 dataset.lookup(id).ok()
289 });
290 Ok(match prepared.execute(dataset)? {
291 QueryResults::Solutions(solutions) => {
292 let variables = solutions.variables().to_vec();
293 let mut rows = Vec::new();
294 for s in solutions {
295 let s = s?;
296 rows.push(variables.iter().map(|v| s.get(v).cloned()).collect());
297 }
298 QueryOutput::Solutions { variables, rows }
299 }
300 QueryResults::Boolean(b) => QueryOutput::Boolean(b),
301 QueryResults::Graph(triples) => QueryOutput::Graph(triples.collect::<Result<Vec<_>, _>>()?),
302 })
303}
304
305fn query_dataset(q: &Query) -> Option<&spargebra::algebra::QueryDataset> {
306 match q {
307 Query::Select { dataset, .. }
308 | Query::Construct { dataset, .. }
309 | Query::Describe { dataset, .. }
310 | Query::Ask { dataset, .. } => dataset.as_ref(),
311 }
312}
313
314pub fn delete_insert<B: SyncBackend>(
317 backend: &B,
318 op: &spargebra::GraphUpdateOperation,
319 base_iri: Option<&oxiri::Iri<String>>,
320) -> Result<(Vec<oxrdf::Quad>, Vec<oxrdf::Quad>)> {
321 let spargebra::GraphUpdateOperation::DeleteInsert {
322 delete,
323 insert,
324 using,
325 pattern,
326 } = op
327 else {
328 return Err(Error::Other("not a DELETE/INSERT operation".into()));
329 };
330 let evaluator = QueryEvaluator::new()
331 .with_custom_function(crate::text::text_match_name(), crate::text::text_match);
332 let prepared = evaluator.prepare_delete_insert(
333 delete.clone(),
334 insert.clone(),
335 base_iri.cloned(),
336 using.clone(),
337 pattern,
338 );
339 let mut deletes = Vec::new();
340 let mut inserts = Vec::new();
341 for q in prepared.execute(SqlDataset::new(backend))? {
342 match q? {
343 spareval::DeleteInsertQuad::Delete(q) => deletes.push(q),
344 spareval::DeleteInsertQuad::Insert(q) => inserts.push(q),
345 }
346 }
347 Ok((deletes, inserts))
348}