1use crate::encoding::{
6 graph_id, named_node_id, subject_id, term_id, EncodedRows, DEFAULT_GRAPH_ID,
7};
8use crate::error::{Error, Result};
9use crate::job::{Job, OneShot, Step};
10use crate::resolve::TermResolver;
11use crate::schema::{create_schema, StoreOptions};
12use crate::sql::{Capabilities, Request, Response, SqlValue, Statement};
13use crate::stats::Stats;
14use crate::writer::{term_statements, EncodedQuads};
15use oxrdf::{
16 GraphName, GraphNameRef, NamedNodeRef, NamedOrBlankNode, NamedOrBlankNodeRef, Quad, QuadRef,
17 Term, TermRef,
18};
19
20fn id_col(caps: &Capabilities, c: &str) -> String {
21 if caps.int64_as_text {
22 format!("CAST({c} AS TEXT)")
23 } else {
24 c.into()
25 }
26}
27
28pub fn open_job(options: &StoreOptions, caps: &Capabilities) -> impl Job<Output = Stats> {
30 let schema = create_schema(options);
31 let load = Stats::load_request(caps);
32 struct Open {
33 schema: Option<Request>,
34 load: Option<Request>,
35 }
36 impl Job for Open {
37 type Output = Stats;
38 fn step(&mut self, response: Option<Response>) -> Result<Step<Stats>> {
39 if let Some(r) = self.schema.take() {
40 return Ok(Step::Execute(r));
41 }
42 if let Some(r) = self.load.take() {
43 return Ok(Step::Execute(r));
44 }
45 Stats::from_response(&response.unwrap_or_default()).map(Step::Done)
46 }
47 }
48 Open {
49 schema: Some(schema),
50 load: Some(load),
51 }
52}
53
54pub fn stats_job(caps: &Capabilities) -> OneShot<Stats> {
56 OneShot::new(Stats::load_request(caps), |r| Stats::from_response(&r))
57}
58
59pub fn optimize_job(caps: &Capabilities) -> impl Job<Output = Stats> {
61 let refresh = Stats::refresh_request();
62 let load = Stats::load_request(caps);
63 struct Optimize {
64 refresh: Option<Request>,
65 load: Option<Request>,
66 }
67 impl Job for Optimize {
68 type Output = Stats;
69 fn step(&mut self, response: Option<Response>) -> Result<Step<Stats>> {
70 if let Some(r) = self.refresh.take() {
71 return Ok(Step::Execute(r));
72 }
73 if let Some(r) = self.load.take() {
74 return Ok(Step::Execute(r));
75 }
76 Stats::from_response(&response.unwrap_or_default()).map(Step::Done)
77 }
78 }
79 Optimize {
80 refresh: Some(refresh),
81 load: Some(load),
82 }
83}
84
85fn count_tail(r: &Response, n: usize) -> u64 {
86 r.iter().rev().take(n).map(|rs| rs.changes).sum()
87}
88
89pub fn insert_job<'a>(
91 quads: impl IntoIterator<Item = QuadRef<'a>>,
92 caps: &Capabilities,
93) -> OneShot<u64> {
94 let quads: Vec<QuadRef<'a>> = quads.into_iter().collect();
95 let schema = quads.iter().any(|q| crate::reason::is_schema_quad(*q));
96 let enc = EncodedQuads::new(quads);
97 let quad_stmts = crate::writer::quad_insert_statements(&enc.quads, caps).len();
98 let mut stmts = enc.insert_statements(caps);
99 let end = stmts.len();
100 if schema {
101 stmts.extend(crate::reason::closure_statements());
102 }
103 OneShot::new(Request::atomic(stmts), move |mut r| {
104 r.truncate(end);
105 Ok(count_tail(&r, quad_stmts))
106 })
107}
108
109pub fn insert_request<'a>(
111 quads: impl IntoIterator<Item = QuadRef<'a>>,
112 caps: &Capabilities,
113) -> Request {
114 Request::atomic(EncodedQuads::new(quads).insert_statements(caps))
115}
116
117pub fn remove_job<'a>(
119 quads: impl IntoIterator<Item = QuadRef<'a>>,
120 caps: &Capabilities,
121) -> OneShot<u64> {
122 let quads: Vec<QuadRef<'a>> = quads.into_iter().collect();
123 let schema = quads.iter().any(|q| crate::reason::is_schema_quad(*q));
124 let enc = EncodedQuads::new(quads);
125 let mut stmts = enc.delete_statements(caps);
126 let end = stmts.len();
127 if schema {
128 stmts.extend(crate::reason::closure_statements());
129 }
130 OneShot::new(Request::atomic(stmts), move |r| {
131 Ok(r.iter().take(end).map(|rs| rs.changes).sum())
132 })
133}
134
135fn scalar(r: &Response) -> i64 {
136 r.first()
137 .and_then(|rs| rs.rows.first())
138 .and_then(|row| row.first())
139 .and_then(SqlValue::as_i64)
140 .unwrap_or(0)
141}
142
143pub fn contains_job(quad: QuadRef<'_>) -> OneShot<bool> {
144 let [s, p, o, g] = [
145 subject_id(quad.subject),
146 named_node_id(quad.predicate.as_str()),
147 term_id(quad.object),
148 graph_id(quad.graph_name),
149 ];
150 OneShot::new(
151 Request::read(vec![Statement::new(format!(
152 "SELECT EXISTS (SELECT 1 FROM quads WHERE s = {s} AND p = {p} AND o = {o} AND g = {g})"
153 ))]),
154 |r| Ok(scalar(&r) != 0),
155 )
156}
157
158pub fn len_job() -> OneShot<usize> {
159 OneShot::new(
160 Request::read(vec!["SELECT COUNT(*) FROM quads".into()]),
161 |r| Ok(scalar(&r) as usize),
162 )
163}
164
165pub fn is_empty_job() -> OneShot<bool> {
166 OneShot::new(
167 Request::read(vec!["SELECT NOT EXISTS (SELECT 1 FROM quads)".into()]),
168 |r| Ok(scalar(&r) != 0),
169 )
170}
171
172pub struct ScanJob {
174 sql: Option<String>,
175 caps: Capabilities,
176 resolver: TermResolver,
177 rows: Vec<[i64; 4]>,
178 started: bool,
179}
180
181impl ScanJob {
182 fn new(where_clause: String, caps: &Capabilities) -> Self {
183 let sql = format!(
184 "SELECT {}, {}, {}, {} FROM quads{}",
185 id_col(caps, "s"),
186 id_col(caps, "p"),
187 id_col(caps, "o"),
188 id_col(caps, "g"),
189 if where_clause.is_empty() {
190 String::new()
191 } else {
192 format!(" WHERE {where_clause}")
193 }
194 );
195 Self {
196 sql: Some(sql),
197 caps: caps.clone(),
198 resolver: TermResolver::default(),
199 rows: Vec::new(),
200 started: false,
201 }
202 }
203
204 fn finish(&self) -> Result<Vec<Quad>> {
205 self.rows
206 .iter()
207 .map(|[s, p, o, g]| {
208 let gname = if *g == DEFAULT_GRAPH_ID {
209 GraphName::DefaultGraph
210 } else {
211 crate::encoding::to_graph_name(*g, Some(self.resolver.get(*g)?))?
212 };
213 crate::encoding::make_quad(
214 self.resolver.get(*s)?,
215 self.resolver.get(*p)?,
216 self.resolver.get(*o)?,
217 gname,
218 )
219 })
220 .collect()
221 }
222}
223
224impl Job for ScanJob {
225 type Output = Vec<Quad>;
226
227 fn step(&mut self, response: Option<Response>) -> Result<Step<Vec<Quad>>> {
228 if let Some(sql) = self.sql.take() {
229 return Ok(Step::Execute(Request::read(vec![Statement::new(sql)])));
230 }
231 let response =
232 response.ok_or_else(|| Error::Other("scan resumed without response".into()))?;
233 if self.started {
234 self.resolver.absorb(response)?;
235 } else {
236 self.started = true;
237 for rs in response {
238 for row in rs.rows {
239 let ids: Vec<i64> = row.iter().filter_map(SqlValue::as_i64).collect();
240 let [s, p, o, g] = ids[..] else {
241 return Err(Error::corrupted("bad quad row"));
242 };
243 for id in [s, p, o, g] {
244 self.resolver.want(id);
245 }
246 self.rows.push([s, p, o, g]);
247 }
248 }
249 }
250 match self.resolver.request(&self.caps) {
251 Some(r) => Ok(Step::Execute(r)),
252 None => self.finish().map(Step::Done),
253 }
254 }
255}
256
257pub fn neighbourhood_job(
260 nodes: &[TermRef<'_>],
261 incoming: bool,
262 predicates: Option<&[NamedNodeRef<'_>]>,
263 default_graph_only: bool,
264 caps: &Capabilities,
265) -> ScanJob {
266 let ids: Vec<String> = nodes.iter().map(|t| term_id(*t).to_string()).collect();
267 let mut w = vec![format!(
268 "{} IN ({})",
269 if incoming { "o" } else { "s" },
270 if ids.is_empty() {
271 "NULL".into()
272 } else {
273 ids.join(",")
274 }
275 )];
276 if let Some(ps) = predicates {
277 let ps: Vec<String> = ps
278 .iter()
279 .map(|p| named_node_id(p.as_str()).to_string())
280 .collect();
281 w.push(format!(
282 "p IN ({})",
283 if ps.is_empty() {
284 "NULL".into()
285 } else {
286 ps.join(",")
287 }
288 ));
289 }
290 if default_graph_only {
291 w.push(format!("g = {DEFAULT_GRAPH_ID}"));
292 }
293 ScanJob::new(w.join(" AND "), caps)
294}
295
296pub fn scan_job(
298 subject: Option<NamedOrBlankNodeRef<'_>>,
299 predicate: Option<NamedNodeRef<'_>>,
300 object: Option<TermRef<'_>>,
301 graph_name: Option<GraphNameRef<'_>>,
302 caps: &Capabilities,
303) -> ScanJob {
304 let mut w = Vec::new();
305 if let Some(s) = subject {
306 w.push(format!("s = {}", subject_id(s)));
307 }
308 if let Some(p) = predicate {
309 w.push(format!("p = {}", named_node_id(p.as_str())));
310 }
311 if let Some(o) = object {
312 w.push(format!("o = {}", term_id(o)));
313 }
314 if let Some(g) = graph_name {
315 w.push(format!("g = {}", graph_id(g)));
316 }
317 ScanJob::new(w.join(" AND "), caps)
318}
319
320pub fn named_graphs_job(caps: &Capabilities) -> impl Job<Output = Vec<NamedOrBlankNode>> {
322 struct Graphs {
323 sql: Option<String>,
324 caps: Capabilities,
325 resolver: TermResolver,
326 ids: Vec<i64>,
327 started: bool,
328 }
329 impl Job for Graphs {
330 type Output = Vec<NamedOrBlankNode>;
331 fn step(&mut self, response: Option<Response>) -> Result<Step<Self::Output>> {
332 if let Some(sql) = self.sql.take() {
333 return Ok(Step::Execute(Request::read(vec![Statement::new(sql)])));
334 }
335 let response = response.unwrap_or_default();
336 if self.started {
337 self.resolver.absorb(response)?;
338 } else {
339 self.started = true;
340 self.ids = crate::resolve::ids_of(&response, 0);
341 for id in &self.ids {
342 self.resolver.want(*id);
343 }
344 }
345 if let Some(r) = self.resolver.request(&self.caps) {
346 return Ok(Step::Execute(r));
347 }
348 self.ids
349 .iter()
350 .map(|id| crate::encoding::to_subject(self.resolver.get(*id)?))
351 .collect::<Result<_>>()
352 .map(Step::Done)
353 }
354 }
355 Graphs {
356 sql: Some(format!("SELECT {} FROM graphs", id_col(caps, "id"))),
357 caps: caps.clone(),
358 resolver: TermResolver::default(),
359 ids: Vec::new(),
360 started: false,
361 }
362}
363
364pub fn contains_named_graph_job(g: NamedOrBlankNodeRef<'_>) -> OneShot<bool> {
365 let id = subject_id(g);
366 OneShot::new(
367 Request::read(vec![Statement::new(format!(
368 "SELECT EXISTS (SELECT 1 FROM graphs WHERE id = {id})"
369 ))]),
370 |r| Ok(scalar(&r) != 0),
371 )
372}
373
374pub fn insert_named_graph_job(g: NamedOrBlankNodeRef<'_>, caps: &Capabilities) -> OneShot<bool> {
376 let mut rows = EncodedRows::default();
377 let id = rows.subject(g);
378 let mut stmts = term_statements(&rows, caps);
379 stmts.push(Statement::new(format!(
380 "INSERT OR IGNORE INTO graphs(id) VALUES ({id})"
381 )));
382 OneShot::new(Request::atomic(stmts), |r| Ok(count_tail(&r, 1) > 0))
383}
384
385pub fn remove_named_graph_job(g: NamedOrBlankNodeRef<'_>) -> OneShot<bool> {
387 let id = subject_id(g);
388 let mut stmts = vec![
389 Statement::new(format!("DELETE FROM quads WHERE g = {id}")),
390 Statement::new(format!("DELETE FROM graphs WHERE id = {id}")),
391 ];
392 stmts.extend(crate::reason::closure_statements());
393 OneShot::new(Request::atomic(stmts), |r| {
394 Ok(r.iter().take(2).any(|rs| rs.changes > 0))
395 })
396}
397
398pub fn clear_graph_job(g: GraphNameRef<'_>) -> OneShot<()> {
400 let id = graph_id(g);
401 let mut stmts = vec![Statement::new(format!("DELETE FROM quads WHERE g = {id}"))];
402 stmts.extend(crate::reason::closure_statements());
403 OneShot::new(Request::atomic(stmts), |_| Ok(()))
404}
405
406pub fn clear_job() -> OneShot<()> {
408 OneShot::new(
409 Request::atomic(vec![
410 "DELETE FROM quads".into(),
411 "DELETE FROM quads_inf".into(),
412 "DELETE FROM tbox_closure".into(),
413 "DELETE FROM graphs".into(),
414 "DELETE FROM triple_terms".into(),
415 "DELETE FROM terms".into(),
416 ]),
417 |_| Ok(()),
418 )
419}
420
421pub fn encode_term(t: &Term) -> i64 {
423 term_id(t.as_ref())
424}
425
426pub fn materialize_job(max_rounds: usize, caps: &Capabilities) -> impl Job<Output = u64> {
431 struct Materialize {
432 reset: Option<Request>,
433 round: usize,
434 max_rounds: usize,
435 counting: bool,
436 }
437 impl Job for Materialize {
438 type Output = u64;
439 fn step(&mut self, response: Option<Response>) -> Result<Step<u64>> {
440 if let Some(r) = self.reset.take() {
441 return Ok(Step::Execute(r));
442 }
443 if self.counting {
444 return Ok(Step::Done(
445 scalar(&response.unwrap_or_default()).max(0) as u64
446 ));
447 }
448 let changed: u64 = response.iter().flatten().map(|rs| rs.changes).sum();
449 if self.round > 0 && (changed == 0 || self.round >= self.max_rounds) {
450 self.counting = true;
451 return Ok(Step::Execute(Request::read(vec![
452 "SELECT COUNT(*) FROM quads_inf".into(),
453 ])));
454 }
455 self.round += 1;
456 Ok(Step::Execute(Request::atomic(
457 crate::reason::materialize_round(),
458 )))
459 }
460 }
461 Materialize {
462 reset: Some(Request::atomic(crate::reason::materialize_reset(caps))),
463 round: 0,
464 max_rounds,
465 counting: false,
466 }
467}
468
469pub fn clear_inferences_job() -> OneShot<()> {
471 OneShot::new(
472 Request::atomic(vec![Statement::new("DELETE FROM quads_inf")]),
473 |_| Ok(()),
474 )
475}