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