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