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
85pub fn schema_refresh_for<'a>(quads: impl IntoIterator<Item = QuadRef<'a>>) -> Vec<Statement> {
92 let mut closure = false;
93 let mut shapes = false;
94 for q in quads {
95 closure |= crate::reason::is_schema_quad(q);
96 shapes |= crate::shapes::is_shape_quad(q);
97 if closure && shapes {
98 break;
99 }
100 }
101 schema_refresh(closure, shapes)
102}
103
104fn schema_refresh(closure: bool, shapes: bool) -> Vec<Statement> {
105 let mut s = Vec::new();
106 if closure {
107 s.extend(crate::reason::closure_statements());
108 }
109 if shapes {
110 s.extend(crate::shapes::refresh_statements());
111 }
112 s
113}
114
115fn count_tail(r: &Response, n: usize) -> u64 {
116 r.iter().rev().take(n).map(|rs| rs.changes).sum()
117}
118
119pub fn insert_job<'a>(
121 quads: impl IntoIterator<Item = QuadRef<'a>>,
122 caps: &Capabilities,
123) -> OneShot<u64> {
124 let quads: Vec<QuadRef<'a>> = quads.into_iter().collect();
125 let refresh = schema_refresh_for(quads.iter().copied());
126 let enc = EncodedQuads::new(quads);
127 let quad_stmts = crate::writer::quad_insert_statements(&enc.quads, caps).len();
128 let mut stmts = enc.insert_statements(caps);
129 let end = stmts.len();
130 stmts.extend(refresh);
131 OneShot::new(Request::atomic(stmts), move |mut r| {
132 r.truncate(end);
133 Ok(count_tail(&r, quad_stmts))
134 })
135}
136
137pub fn insert_request<'a>(
139 quads: impl IntoIterator<Item = QuadRef<'a>>,
140 caps: &Capabilities,
141) -> Request {
142 Request::atomic(EncodedQuads::new(quads).insert_statements(caps))
143}
144
145pub fn remove_job<'a>(
147 quads: impl IntoIterator<Item = QuadRef<'a>>,
148 caps: &Capabilities,
149) -> OneShot<u64> {
150 let quads: Vec<QuadRef<'a>> = quads.into_iter().collect();
151 let refresh = schema_refresh_for(quads.iter().copied());
152 let enc = EncodedQuads::new(quads);
153 let mut stmts = enc.delete_statements(caps);
154 let end = stmts.len();
155 stmts.extend(refresh);
156 OneShot::new(Request::atomic(stmts), move |r| {
157 Ok(r.iter().take(end).map(|rs| rs.changes).sum())
158 })
159}
160
161fn scalar(r: &Response) -> i64 {
162 r.first()
163 .and_then(|rs| rs.rows.first())
164 .and_then(|row| row.first())
165 .and_then(SqlValue::as_i64)
166 .unwrap_or(0)
167}
168
169pub fn contains_job(quad: QuadRef<'_>) -> OneShot<bool> {
170 let [s, p, o, g] = [
171 subject_id(quad.subject),
172 named_node_id(quad.predicate.as_str()),
173 term_id(quad.object),
174 graph_id(quad.graph_name),
175 ];
176 OneShot::new(
177 Request::read(vec![Statement::new(format!(
178 "SELECT EXISTS (SELECT 1 FROM quads WHERE s = {s} AND p = {p} AND o = {o} AND g = {g})"
179 ))]),
180 |r| Ok(scalar(&r) != 0),
181 )
182}
183
184pub fn len_job() -> OneShot<usize> {
185 OneShot::new(
186 Request::read(vec!["SELECT COUNT(*) FROM quads".into()]),
187 |r| Ok(scalar(&r) as usize),
188 )
189}
190
191pub fn is_empty_job() -> OneShot<bool> {
192 OneShot::new(
193 Request::read(vec!["SELECT NOT EXISTS (SELECT 1 FROM quads)".into()]),
194 |r| Ok(scalar(&r) != 0),
195 )
196}
197
198pub struct ScanJob {
200 sql: Option<String>,
201 caps: Capabilities,
202 resolver: TermResolver,
203 rows: Vec<[i64; 4]>,
204 started: bool,
205}
206
207impl ScanJob {
208 fn new(where_clause: String, caps: &Capabilities) -> Self {
209 let sql = format!(
210 "SELECT {}, {}, {}, {} FROM quads{}",
211 id_col(caps, "s"),
212 id_col(caps, "p"),
213 id_col(caps, "o"),
214 id_col(caps, "g"),
215 if where_clause.is_empty() {
216 String::new()
217 } else {
218 format!(" WHERE {where_clause}")
219 }
220 );
221 Self {
222 sql: Some(sql),
223 caps: caps.clone(),
224 resolver: TermResolver::default(),
225 rows: Vec::new(),
226 started: false,
227 }
228 }
229
230 fn finish(&self) -> Result<Vec<Quad>> {
231 self.rows
232 .iter()
233 .map(|[s, p, o, g]| {
234 let gname = if *g == DEFAULT_GRAPH_ID {
235 GraphName::DefaultGraph
236 } else {
237 crate::encoding::to_graph_name(*g, Some(self.resolver.get(*g)?))?
238 };
239 crate::encoding::make_quad(
240 self.resolver.get(*s)?,
241 self.resolver.get(*p)?,
242 self.resolver.get(*o)?,
243 gname,
244 )
245 })
246 .collect()
247 }
248}
249
250impl Job for ScanJob {
251 type Output = Vec<Quad>;
252
253 fn step(&mut self, response: Option<Response>) -> Result<Step<Vec<Quad>>> {
254 if let Some(sql) = self.sql.take() {
255 return Ok(Step::Execute(Request::read(vec![Statement::new(sql)])));
256 }
257 let response =
258 response.ok_or_else(|| Error::Other("scan resumed without response".into()))?;
259 if self.started {
260 self.resolver.absorb(response)?;
261 } else {
262 self.started = true;
263 for rs in response {
264 for row in rs.rows {
265 let ids: Vec<i64> = row.iter().filter_map(SqlValue::as_i64).collect();
266 let [s, p, o, g] = ids[..] else {
267 return Err(Error::corrupted("bad quad row"));
268 };
269 for id in [s, p, o, g] {
270 self.resolver.want(id);
271 }
272 self.rows.push([s, p, o, g]);
273 }
274 }
275 }
276 match self.resolver.request(&self.caps) {
277 Some(r) => Ok(Step::Execute(r)),
278 None => self.finish().map(Step::Done),
279 }
280 }
281}
282
283pub fn neighbourhood_job(
286 nodes: &[TermRef<'_>],
287 incoming: bool,
288 predicates: Option<&[NamedNodeRef<'_>]>,
289 default_graph_only: bool,
290 caps: &Capabilities,
291) -> ScanJob {
292 let ids: Vec<String> = nodes.iter().map(|t| term_id(*t).to_string()).collect();
293 let mut w = vec![format!(
294 "{} IN ({})",
295 if incoming { "o" } else { "s" },
296 if ids.is_empty() {
297 "NULL".into()
298 } else {
299 ids.join(",")
300 }
301 )];
302 if let Some(ps) = predicates {
303 let ps: Vec<String> = ps
304 .iter()
305 .map(|p| named_node_id(p.as_str()).to_string())
306 .collect();
307 w.push(format!(
308 "p IN ({})",
309 if ps.is_empty() {
310 "NULL".into()
311 } else {
312 ps.join(",")
313 }
314 ));
315 }
316 if default_graph_only {
317 w.push(format!("g = {DEFAULT_GRAPH_ID}"));
318 }
319 ScanJob::new(w.join(" AND "), caps)
320}
321
322pub fn scan_job(
324 subject: Option<NamedOrBlankNodeRef<'_>>,
325 predicate: Option<NamedNodeRef<'_>>,
326 object: Option<TermRef<'_>>,
327 graph_name: Option<GraphNameRef<'_>>,
328 caps: &Capabilities,
329) -> ScanJob {
330 let mut w = Vec::new();
331 if let Some(s) = subject {
332 w.push(format!("s = {}", subject_id(s)));
333 }
334 if let Some(p) = predicate {
335 w.push(format!("p = {}", named_node_id(p.as_str())));
336 }
337 if let Some(o) = object {
338 w.push(format!("o = {}", term_id(o)));
339 }
340 if let Some(g) = graph_name {
341 w.push(format!("g = {}", graph_id(g)));
342 }
343 ScanJob::new(w.join(" AND "), caps)
344}
345
346pub fn named_graphs_job(caps: &Capabilities) -> impl Job<Output = Vec<NamedOrBlankNode>> {
348 struct Graphs {
349 sql: Option<String>,
350 caps: Capabilities,
351 resolver: TermResolver,
352 ids: Vec<i64>,
353 started: bool,
354 }
355 impl Job for Graphs {
356 type Output = Vec<NamedOrBlankNode>;
357 fn step(&mut self, response: Option<Response>) -> Result<Step<Self::Output>> {
358 if let Some(sql) = self.sql.take() {
359 return Ok(Step::Execute(Request::read(vec![Statement::new(sql)])));
360 }
361 let response = response.unwrap_or_default();
362 if self.started {
363 self.resolver.absorb(response)?;
364 } else {
365 self.started = true;
366 self.ids = crate::resolve::ids_of(&response, 0);
367 for id in &self.ids {
368 self.resolver.want(*id);
369 }
370 }
371 if let Some(r) = self.resolver.request(&self.caps) {
372 return Ok(Step::Execute(r));
373 }
374 self.ids
375 .iter()
376 .map(|id| crate::encoding::to_subject(self.resolver.get(*id)?))
377 .collect::<Result<_>>()
378 .map(Step::Done)
379 }
380 }
381 Graphs {
382 sql: Some(format!("SELECT {} FROM graphs", id_col(caps, "id"))),
383 caps: caps.clone(),
384 resolver: TermResolver::default(),
385 ids: Vec::new(),
386 started: false,
387 }
388}
389
390pub fn contains_named_graph_job(g: NamedOrBlankNodeRef<'_>) -> OneShot<bool> {
391 let id = subject_id(g);
392 OneShot::new(
393 Request::read(vec![Statement::new(format!(
394 "SELECT EXISTS (SELECT 1 FROM graphs WHERE id = {id})"
395 ))]),
396 |r| Ok(scalar(&r) != 0),
397 )
398}
399
400pub fn insert_named_graph_job(g: NamedOrBlankNodeRef<'_>, caps: &Capabilities) -> OneShot<bool> {
402 let mut rows = EncodedRows::default();
403 let id = rows.subject(g);
404 let mut stmts = term_statements(&rows, caps);
405 stmts.push(Statement::new(format!(
406 "INSERT OR IGNORE INTO graphs(id) VALUES ({id})"
407 )));
408 OneShot::new(Request::atomic(stmts), |r| Ok(count_tail(&r, 1) > 0))
409}
410
411pub fn remove_named_graph_job(g: NamedOrBlankNodeRef<'_>) -> OneShot<bool> {
413 let id = subject_id(g);
414 let mut stmts = vec![
415 Statement::new(format!("DELETE FROM quads WHERE g = {id}")),
416 Statement::new(format!("DELETE FROM graphs WHERE id = {id}")),
417 ];
418 stmts.extend(crate::registry::unregister_statements(id));
419 stmts.extend(schema_refresh(true, true));
420 OneShot::new(Request::atomic(stmts), |r| {
421 Ok(r.iter().take(2).any(|rs| rs.changes > 0))
422 })
423}
424
425pub fn clear_graph_job(g: GraphNameRef<'_>) -> OneShot<()> {
427 let id = graph_id(g);
428 let mut stmts = vec![Statement::new(format!("DELETE FROM quads WHERE g = {id}"))];
429 stmts.extend(schema_refresh(true, true));
430 OneShot::new(Request::atomic(stmts), |_| Ok(()))
431}
432
433pub fn clear_job() -> OneShot<()> {
435 OneShot::new(
436 Request::atomic(vec![
437 "DELETE FROM quads".into(),
438 "DELETE FROM quads_inf".into(),
439 "DELETE FROM quads_inf_src".into(),
440 "DELETE FROM inf_producers".into(),
441 "DELETE FROM tbox_closure".into(),
442 "DELETE FROM shapes_index".into(),
443 "DELETE FROM shapes_in".into(),
444 "DELETE FROM schema_graphs".into(),
445 "DELETE FROM graphs".into(),
446 "DELETE FROM triple_terms".into(),
447 "DELETE FROM terms".into(),
448 ]),
449 |_| Ok(()),
450 )
451}
452
453pub fn register_schema_graph_job(
462 entry: &crate::registry::SchemaGraph,
463 graph: GraphNameRef<'_>,
464 caps: &Capabilities,
465) -> OneShot<()> {
466 let mut stmts = Vec::new();
467 let named: Option<NamedOrBlankNodeRef<'_>> = match graph {
468 GraphNameRef::NamedNode(n) => Some(n.into()),
469 GraphNameRef::BlankNode(b) => Some(b.into()),
470 GraphNameRef::DefaultGraph => None,
471 };
472 if let Some(g) = named {
473 let mut rows = EncodedRows::default();
474 let id = rows.subject(g);
475 stmts.extend(term_statements(&rows, caps));
476 stmts.push(Statement::new(format!(
477 "INSERT OR IGNORE INTO graphs(id) VALUES ({id})"
478 )));
479 }
480 stmts.extend(crate::registry::register_statements(entry));
481 stmts.extend(schema_refresh(true, true));
482 OneShot::new(Request::atomic(stmts), |_| Ok(()))
483}
484
485pub fn unregister_schema_graph_job(graph: i64) -> OneShot<bool> {
487 let mut stmts = crate::registry::unregister_statements(graph);
488 stmts.extend(schema_refresh(true, true));
489 OneShot::new(Request::atomic(stmts), |r| Ok(scalar_changes(&r, 1) > 0))
490}
491
492pub fn set_schema_graph_active_job(graph: i64, active: bool) -> OneShot<bool> {
494 let mut stmts = crate::registry::set_active_statements(graph, active);
495 stmts.extend(schema_refresh(true, true));
496 OneShot::new(Request::atomic(stmts), |r| Ok(scalar_changes(&r, 1) > 0))
497}
498
499pub fn drop_schema_graph_job(graph: i64) -> OneShot<u64> {
501 let mut stmts = crate::registry::drop_statements(graph);
502 stmts.extend(schema_refresh(true, true));
503 OneShot::new(Request::atomic(stmts), |r| Ok(scalar_changes(&r, 1)))
504}
505
506pub fn schema_graphs_job(
508 caps: &Capabilities,
509) -> impl Job<Output = Vec<(crate::registry::SchemaGraph, GraphName)>> {
510 struct Registry {
511 request: Option<Request>,
512 caps: Capabilities,
513 resolver: TermResolver,
514 rows: Vec<crate::registry::SchemaGraph>,
515 started: bool,
516 }
517 impl Job for Registry {
518 type Output = Vec<(crate::registry::SchemaGraph, GraphName)>;
519 fn step(&mut self, response: Option<Response>) -> Result<Step<Self::Output>> {
520 if let Some(r) = self.request.take() {
521 return Ok(Step::Execute(r));
522 }
523 let response = response.unwrap_or_default();
524 if self.started {
525 self.resolver.absorb(response)?;
526 } else {
527 self.started = true;
528 self.rows = crate::registry::from_response(&response)?;
529 for row in &self.rows {
530 if row.graph != DEFAULT_GRAPH_ID {
531 self.resolver.want(row.graph);
532 }
533 }
534 }
535 if let Some(r) = self.resolver.request(&self.caps) {
536 return Ok(Step::Execute(r));
537 }
538 std::mem::take(&mut self.rows)
539 .into_iter()
540 .map(|row| {
541 let name = if row.graph == DEFAULT_GRAPH_ID {
542 GraphName::DefaultGraph
543 } else {
544 crate::encoding::to_graph_name(
545 row.graph,
546 Some(self.resolver.get(row.graph)?),
547 )?
548 };
549 Ok((row, name))
550 })
551 .collect::<Result<Vec<_>>>()
552 .map(Step::Done)
553 }
554 }
555 Registry {
556 request: Some(crate::registry::load_request(caps)),
557 caps: caps.clone(),
558 resolver: TermResolver::default(),
559 rows: Vec::new(),
560 started: false,
561 }
562}
563
564pub fn shape_index_job(caps: &Capabilities) -> OneShot<crate::shapes::ShapeIndex> {
566 OneShot::new(crate::shapes::ShapeIndex::load_request(caps), |r| {
567 crate::shapes::ShapeIndex::from_response(&r)
568 })
569}
570
571fn scalar_changes(r: &Response, n: usize) -> u64 {
573 r.iter().take(n).map(|rs| rs.changes).sum()
574}
575
576pub fn encode_term(t: &Term) -> i64 {
578 term_id(t.as_ref())
579}
580
581pub fn materialize_job(max_rounds: usize, caps: &Capabilities) -> impl Job<Output = u64> {
586 struct Materialize {
587 reset: Option<Request>,
588 round: usize,
589 max_rounds: usize,
590 counting: bool,
591 }
592 impl Job for Materialize {
593 type Output = u64;
594 fn step(&mut self, response: Option<Response>) -> Result<Step<u64>> {
595 if let Some(r) = self.reset.take() {
596 return Ok(Step::Execute(r));
597 }
598 if self.counting {
599 return Ok(Step::Done(
600 scalar(&response.unwrap_or_default()).max(0) as u64
601 ));
602 }
603 let changed: u64 = response.iter().flatten().map(|rs| rs.changes).sum();
604 if self.round > 0 && (changed == 0 || self.round >= self.max_rounds) {
605 self.counting = true;
606 return Ok(Step::Execute(Request::read(vec![
607 "SELECT COUNT(*) FROM quads_inf".into(),
608 ])));
609 }
610 self.round += 1;
611 Ok(Step::Execute(Request::atomic(
612 crate::reason::materialize_round(),
613 )))
614 }
615 }
616 Materialize {
617 reset: Some(Request::atomic(crate::reason::materialize_reset(caps))),
618 round: 0,
619 max_rounds,
620 counting: false,
621 }
622}
623
624pub fn clear_inferences_job() -> OneShot<()> {
626 OneShot::new(
627 Request::atomic(vec![
628 Statement::new("DELETE FROM quads_inf"),
629 Statement::new("DELETE FROM quads_inf_src"),
630 ]),
631 |_| Ok(()),
632 )
633}
634
635pub fn clear_inferences_of_job(producer: &str) -> OneShot<()> {
637 OneShot::new(
638 Request::atomic(crate::reason::inference_reset(producer)),
639 |_| Ok(()),
640 )
641}
642
643pub fn inference_producers_job(quad: QuadRef<'_>) -> OneShot<Vec<String>> {
646 let [s, p, o, g] = [
647 subject_id(quad.subject),
648 named_node_id(quad.predicate.as_str()),
649 term_id(quad.object),
650 graph_id(quad.graph_name),
651 ];
652 OneShot::new(
653 Request::read(vec![Statement::new(format!(
654 "SELECT p.name FROM quads_inf_src q JOIN inf_producers p ON p.id = q.src \
655 WHERE q.s = {s} AND q.p = {p} AND q.o = {o} AND q.g = {g} ORDER BY p.name"
656 ))]),
657 |r| {
658 Ok(r.first()
659 .map(|rs| {
660 rs.rows
661 .iter()
662 .filter_map(|row| match row.first() {
663 Some(SqlValue::Text(t)) => Some(t.clone()),
664 _ => None,
665 })
666 .collect()
667 })
668 .unwrap_or_default())
669 },
670 )
671}