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::{base_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> {
34 enum Phase {
35 Start,
36 Schema,
37 Stats,
38 Empty(Stats),
39 Upgraded,
40 Reloaded,
41 }
42 struct Open {
43 options: StoreOptions,
44 caps: Capabilities,
45 phase: Phase,
46 }
47 impl Job for Open {
48 type Output = Stats;
49 fn step(&mut self, response: Option<Response>) -> Result<Step<Stats>> {
50 let wanted = self.options.versioning;
51 match std::mem::replace(&mut self.phase, Phase::Reloaded) {
52 Phase::Start => {
53 self.phase = Phase::Schema;
54 Ok(Step::Execute(base_schema(&self.options)))
55 }
56 Phase::Schema => {
57 self.phase = Phase::Stats;
58 Ok(Step::Execute(Stats::load_request(&self.caps)))
59 }
60 Phase::Stats => {
61 let stats = Stats::from_response(&response.unwrap_or_default())?;
62 if wanted <= stats.version.level {
63 return Ok(Step::Done(stats));
64 }
65 self.phase = Phase::Empty(stats);
66 Ok(Step::Execute(Request::read(vec![Statement::new(
67 "SELECT NOT EXISTS (SELECT 1 FROM quads)",
68 )])))
69 }
70 Phase::Empty(stats) => {
71 let response = response.unwrap_or_default();
72 let empty = response
73 .first()
74 .and_then(|r| r.rows.first())
75 .and_then(|row| row.first())
76 .and_then(SqlValue::as_i64)
77 == Some(1);
78 if !empty || stats.version.history != crate::version::History::None {
79 return Err(Error::Other(format!(
80 "the store's versioning level is `{}`; opening does not change it. Raise it explicitly (`Store::set_versioning`, `oxilite versioning set {wanted}`)",
81 stats.version.level
82 )));
83 }
84 self.phase = Phase::Upgraded;
85 Ok(Step::Execute(Request::atomic(
86 crate::version::change_statements(
87 &stats.version,
88 wanted,
89 &self.options.level_change(),
90 )?,
91 )))
92 }
93 Phase::Upgraded => {
94 self.phase = Phase::Reloaded;
95 Ok(Step::Execute(Stats::load_request(&self.caps)))
96 }
97 Phase::Reloaded => {
98 Stats::from_response(&response.unwrap_or_default()).map(Step::Done)
99 }
100 }
101 }
102 }
103 Open {
104 options: options.clone(),
105 caps: caps.clone(),
106 phase: Phase::Start,
107 }
108}
109
110pub fn stats_job(caps: &Capabilities) -> OneShot<Stats> {
112 OneShot::new(Stats::load_request(caps), |r| Stats::from_response(&r))
113}
114
115pub fn optimize_job(caps: &Capabilities) -> impl Job<Output = Stats> {
117 let refresh = Stats::refresh_request();
118 let load = Stats::load_request(caps);
119 struct Optimize {
120 refresh: Option<Request>,
121 load: Option<Request>,
122 }
123 impl Job for Optimize {
124 type Output = Stats;
125 fn step(&mut self, response: Option<Response>) -> Result<Step<Stats>> {
126 if let Some(r) = self.refresh.take() {
127 return Ok(Step::Execute(r));
128 }
129 if let Some(r) = self.load.take() {
130 return Ok(Step::Execute(r));
131 }
132 Stats::from_response(&response.unwrap_or_default()).map(Step::Done)
133 }
134 }
135 Optimize {
136 refresh: Some(refresh),
137 load: Some(load),
138 }
139}
140
141pub fn schema_refresh_for<'a>(quads: impl IntoIterator<Item = QuadRef<'a>>) -> Vec<Statement> {
148 let mut closure = false;
149 let mut shapes = false;
150 for q in quads {
151 closure |= crate::reason::is_schema_quad(q);
152 shapes |= crate::shapes::is_shape_quad(q);
153 if closure && shapes {
154 break;
155 }
156 }
157 schema_refresh(closure, shapes)
158}
159
160fn schema_refresh(closure: bool, shapes: bool) -> Vec<Statement> {
161 let mut s = Vec::new();
162 if closure {
163 s.extend(crate::reason::closure_statements());
164 }
165 if shapes {
166 s.extend(crate::shapes::refresh_statements());
167 }
168 s
169}
170
171fn count_tail(r: &Response, n: usize) -> u64 {
172 r.iter().rev().take(n).map(|rs| rs.changes).sum()
173}
174
175pub fn insert_job<'a>(
177 quads: impl IntoIterator<Item = QuadRef<'a>>,
178 caps: &Capabilities,
179) -> OneShot<u64> {
180 let quads: Vec<QuadRef<'a>> = quads.into_iter().collect();
181 let refresh = schema_refresh_for(quads.iter().copied());
182 let enc = EncodedQuads::new(quads);
183 let quad_stmts = crate::writer::quad_insert_statements(&enc.quads, caps).len();
184 let mut stmts = enc.insert_statements(caps);
185 let end = stmts.len();
186 stmts.extend(refresh);
187 OneShot::new(Request::atomic(stmts), move |mut r| {
188 r.truncate(end);
189 Ok(count_tail(&r, quad_stmts))
190 })
191}
192
193pub fn insert_request<'a>(
195 quads: impl IntoIterator<Item = QuadRef<'a>>,
196 caps: &Capabilities,
197) -> Request {
198 Request::atomic(EncodedQuads::new(quads).insert_statements(caps))
199}
200
201pub fn remove_job<'a>(
203 quads: impl IntoIterator<Item = QuadRef<'a>>,
204 caps: &Capabilities,
205) -> OneShot<u64> {
206 let quads: Vec<QuadRef<'a>> = quads.into_iter().collect();
207 let refresh = schema_refresh_for(quads.iter().copied());
208 let enc = EncodedQuads::new(quads);
209 let mut stmts = enc.delete_statements(caps);
210 let end = stmts.len();
211 stmts.extend(refresh);
212 OneShot::new(Request::atomic(stmts), move |r| {
213 Ok(r.iter().take(end).map(|rs| rs.changes).sum())
214 })
215}
216
217fn scalar(r: &Response) -> i64 {
218 r.first()
219 .and_then(|rs| rs.rows.first())
220 .and_then(|row| row.first())
221 .and_then(SqlValue::as_i64)
222 .unwrap_or(0)
223}
224
225pub fn contains_job(quad: QuadRef<'_>) -> OneShot<bool> {
226 let [s, p, o, g] = [
227 subject_id(quad.subject),
228 named_node_id(quad.predicate.as_str()),
229 term_id(quad.object),
230 graph_id(quad.graph_name),
231 ];
232 OneShot::new(
233 Request::read(vec![Statement::new(format!(
234 "SELECT EXISTS (SELECT 1 FROM quads WHERE s = {s} AND p = {p} AND o = {o} AND g = {g})"
235 ))]),
236 |r| Ok(scalar(&r) != 0),
237 )
238}
239
240pub fn len_job() -> OneShot<usize> {
241 OneShot::new(
242 Request::read(vec!["SELECT COUNT(*) FROM quads".into()]),
243 |r| Ok(scalar(&r) as usize),
244 )
245}
246
247pub fn is_empty_job() -> OneShot<bool> {
248 OneShot::new(
249 Request::read(vec!["SELECT NOT EXISTS (SELECT 1 FROM quads)".into()]),
250 |r| Ok(scalar(&r) != 0),
251 )
252}
253
254pub struct ScanJob {
256 sql: Option<String>,
257 caps: Capabilities,
258 resolver: TermResolver,
259 rows: Vec<[i64; 4]>,
260 started: bool,
261}
262
263impl ScanJob {
264 fn new(where_clause: String, caps: &Capabilities) -> Self {
265 let sql = format!(
266 "SELECT {}, {}, {}, {} FROM quads{}",
267 id_col(caps, "s"),
268 id_col(caps, "p"),
269 id_col(caps, "o"),
270 id_col(caps, "g"),
271 if where_clause.is_empty() {
272 String::new()
273 } else {
274 format!(" WHERE {where_clause}")
275 }
276 );
277 Self {
278 sql: Some(sql),
279 caps: caps.clone(),
280 resolver: TermResolver::default(),
281 rows: Vec::new(),
282 started: false,
283 }
284 }
285
286 fn finish(&self) -> Result<Vec<Quad>> {
287 self.rows
288 .iter()
289 .map(|[s, p, o, g]| {
290 let gname = if *g == DEFAULT_GRAPH_ID {
291 GraphName::DefaultGraph
292 } else {
293 crate::encoding::to_graph_name(*g, Some(self.resolver.get(*g)?))?
294 };
295 crate::encoding::make_quad(
296 self.resolver.get(*s)?,
297 self.resolver.get(*p)?,
298 self.resolver.get(*o)?,
299 gname,
300 )
301 })
302 .collect()
303 }
304}
305
306impl Job for ScanJob {
307 type Output = Vec<Quad>;
308
309 fn step(&mut self, response: Option<Response>) -> Result<Step<Vec<Quad>>> {
310 if let Some(sql) = self.sql.take() {
311 return Ok(Step::Execute(Request::read(vec![Statement::new(sql)])));
312 }
313 let response =
314 response.ok_or_else(|| Error::Other("scan resumed without response".into()))?;
315 if self.started {
316 self.resolver.absorb(response)?;
317 } else {
318 self.started = true;
319 for rs in response {
320 for row in rs.rows {
321 let ids: Vec<i64> = row.iter().filter_map(SqlValue::as_i64).collect();
322 let [s, p, o, g] = ids[..] else {
323 return Err(Error::corrupted("bad quad row"));
324 };
325 for id in [s, p, o, g] {
326 self.resolver.want(id);
327 }
328 self.rows.push([s, p, o, g]);
329 }
330 }
331 }
332 match self.resolver.request(&self.caps) {
333 Some(r) => Ok(Step::Execute(r)),
334 None => self.finish().map(Step::Done),
335 }
336 }
337}
338
339pub fn neighbourhood_job(
342 nodes: &[TermRef<'_>],
343 incoming: bool,
344 predicates: Option<&[NamedNodeRef<'_>]>,
345 default_graph_only: bool,
346 caps: &Capabilities,
347) -> ScanJob {
348 let ids: Vec<String> = nodes.iter().map(|t| term_id(*t).to_string()).collect();
349 let mut w = vec![format!(
350 "{} IN ({})",
351 if incoming { "o" } else { "s" },
352 if ids.is_empty() {
353 "NULL".into()
354 } else {
355 ids.join(",")
356 }
357 )];
358 if let Some(ps) = predicates {
359 let ps: Vec<String> = ps
360 .iter()
361 .map(|p| named_node_id(p.as_str()).to_string())
362 .collect();
363 w.push(format!(
364 "p IN ({})",
365 if ps.is_empty() {
366 "NULL".into()
367 } else {
368 ps.join(",")
369 }
370 ));
371 }
372 if default_graph_only {
373 w.push(format!("g = {DEFAULT_GRAPH_ID}"));
374 }
375 ScanJob::new(w.join(" AND "), caps)
376}
377
378pub fn scan_job(
380 subject: Option<NamedOrBlankNodeRef<'_>>,
381 predicate: Option<NamedNodeRef<'_>>,
382 object: Option<TermRef<'_>>,
383 graph_name: Option<GraphNameRef<'_>>,
384 caps: &Capabilities,
385) -> ScanJob {
386 let mut w = Vec::new();
387 if let Some(s) = subject {
388 w.push(format!("s = {}", subject_id(s)));
389 }
390 if let Some(p) = predicate {
391 w.push(format!("p = {}", named_node_id(p.as_str())));
392 }
393 if let Some(o) = object {
394 w.push(format!("o = {}", term_id(o)));
395 }
396 if let Some(g) = graph_name {
397 w.push(format!("g = {}", graph_id(g)));
398 }
399 ScanJob::new(w.join(" AND "), caps)
400}
401
402pub fn named_graphs_job(caps: &Capabilities) -> impl Job<Output = Vec<NamedOrBlankNode>> {
404 struct Graphs {
405 sql: Option<String>,
406 caps: Capabilities,
407 resolver: TermResolver,
408 ids: Vec<i64>,
409 started: bool,
410 }
411 impl Job for Graphs {
412 type Output = Vec<NamedOrBlankNode>;
413 fn step(&mut self, response: Option<Response>) -> Result<Step<Self::Output>> {
414 if let Some(sql) = self.sql.take() {
415 return Ok(Step::Execute(Request::read(vec![Statement::new(sql)])));
416 }
417 let response = response.unwrap_or_default();
418 if self.started {
419 self.resolver.absorb(response)?;
420 } else {
421 self.started = true;
422 self.ids = crate::resolve::ids_of(&response, 0);
423 for id in &self.ids {
424 self.resolver.want(*id);
425 }
426 }
427 if let Some(r) = self.resolver.request(&self.caps) {
428 return Ok(Step::Execute(r));
429 }
430 self.ids
431 .iter()
432 .map(|id| crate::encoding::to_subject(self.resolver.get(*id)?))
433 .collect::<Result<_>>()
434 .map(Step::Done)
435 }
436 }
437 Graphs {
438 sql: Some(format!("SELECT {} FROM graphs", id_col(caps, "id"))),
439 caps: caps.clone(),
440 resolver: TermResolver::default(),
441 ids: Vec::new(),
442 started: false,
443 }
444}
445
446pub fn contains_named_graph_job(g: NamedOrBlankNodeRef<'_>) -> OneShot<bool> {
447 let id = subject_id(g);
448 OneShot::new(
449 Request::read(vec![Statement::new(format!(
450 "SELECT EXISTS (SELECT 1 FROM graphs WHERE id = {id})"
451 ))]),
452 |r| Ok(scalar(&r) != 0),
453 )
454}
455
456pub fn insert_named_graph_job(g: NamedOrBlankNodeRef<'_>, caps: &Capabilities) -> OneShot<bool> {
458 let mut rows = EncodedRows::default();
459 let id = rows.subject(g);
460 let mut stmts = term_statements(&rows, caps);
461 stmts.push(Statement::new(format!(
462 "INSERT OR IGNORE INTO graphs(id) VALUES ({id})"
463 )));
464 OneShot::new(Request::atomic(stmts), |r| Ok(count_tail(&r, 1) > 0))
465}
466
467pub fn remove_named_graph_job(g: NamedOrBlankNodeRef<'_>) -> OneShot<bool> {
469 let id = subject_id(g);
470 let mut stmts = vec![
471 Statement::new(format!("DELETE FROM quads WHERE g = {id}")),
472 Statement::new(format!("DELETE FROM graphs WHERE id = {id}")),
473 ];
474 stmts.extend(crate::registry::unregister_statements(id));
475 stmts.extend(schema_refresh(true, true));
476 OneShot::new(Request::atomic(stmts), |r| {
477 Ok(r.iter().take(2).any(|rs| rs.changes > 0))
478 })
479}
480
481pub fn clear_graph_job(g: GraphNameRef<'_>) -> OneShot<()> {
483 let id = graph_id(g);
484 let mut stmts = vec![Statement::new(format!("DELETE FROM quads WHERE g = {id}"))];
485 stmts.extend(schema_refresh(true, true));
486 OneShot::new(Request::atomic(stmts), |_| Ok(()))
487}
488
489pub fn clear_job() -> OneShot<()> {
491 OneShot::new(
492 Request::atomic(vec![
493 "DELETE FROM quads".into(),
494 "DELETE FROM quads_inf".into(),
495 "DELETE FROM quads_inf_src".into(),
496 "DELETE FROM inf_producers".into(),
497 "DELETE FROM tbox_closure".into(),
498 "DELETE FROM shapes_index".into(),
499 "DELETE FROM shapes_in".into(),
500 "DELETE FROM schema_graphs".into(),
501 "DELETE FROM graphs".into(),
502 "DELETE FROM triple_terms".into(),
503 "DELETE FROM terms".into(),
504 ]),
505 |_| Ok(()),
506 )
507}
508
509pub fn register_schema_graph_job(
518 entry: &crate::registry::SchemaGraph,
519 graph: GraphNameRef<'_>,
520 caps: &Capabilities,
521) -> OneShot<()> {
522 let mut stmts = Vec::new();
523 let named: Option<NamedOrBlankNodeRef<'_>> = match graph {
524 GraphNameRef::NamedNode(n) => Some(n.into()),
525 GraphNameRef::BlankNode(b) => Some(b.into()),
526 GraphNameRef::DefaultGraph => None,
527 };
528 if let Some(g) = named {
529 let mut rows = EncodedRows::default();
530 let id = rows.subject(g);
531 stmts.extend(term_statements(&rows, caps));
532 stmts.push(Statement::new(format!(
533 "INSERT OR IGNORE INTO graphs(id) VALUES ({id})"
534 )));
535 }
536 stmts.extend(crate::registry::register_statements(entry));
537 stmts.extend(schema_refresh(true, true));
538 OneShot::new(Request::atomic(stmts), |_| Ok(()))
539}
540
541pub fn unregister_schema_graph_job(graph: i64) -> OneShot<bool> {
543 let mut stmts = crate::registry::unregister_statements(graph);
544 stmts.extend(schema_refresh(true, true));
545 OneShot::new(Request::atomic(stmts), |r| Ok(scalar_changes(&r, 1) > 0))
546}
547
548pub fn set_schema_graph_active_job(graph: i64, active: bool) -> OneShot<bool> {
550 let mut stmts = crate::registry::set_active_statements(graph, active);
551 stmts.extend(schema_refresh(true, true));
552 OneShot::new(Request::atomic(stmts), |r| Ok(scalar_changes(&r, 1) > 0))
553}
554
555pub fn drop_schema_graph_job(graph: i64) -> OneShot<u64> {
557 let mut stmts = crate::registry::drop_statements(graph);
558 stmts.extend(schema_refresh(true, true));
559 OneShot::new(Request::atomic(stmts), |r| Ok(scalar_changes(&r, 1)))
560}
561
562pub fn schema_graphs_job(
564 caps: &Capabilities,
565) -> impl Job<Output = Vec<(crate::registry::SchemaGraph, GraphName)>> {
566 struct Registry {
567 request: Option<Request>,
568 caps: Capabilities,
569 resolver: TermResolver,
570 rows: Vec<crate::registry::SchemaGraph>,
571 started: bool,
572 }
573 impl Job for Registry {
574 type Output = Vec<(crate::registry::SchemaGraph, GraphName)>;
575 fn step(&mut self, response: Option<Response>) -> Result<Step<Self::Output>> {
576 if let Some(r) = self.request.take() {
577 return Ok(Step::Execute(r));
578 }
579 let response = response.unwrap_or_default();
580 if self.started {
581 self.resolver.absorb(response)?;
582 } else {
583 self.started = true;
584 self.rows = crate::registry::from_response(&response)?;
585 for row in &self.rows {
586 if row.graph != DEFAULT_GRAPH_ID {
587 self.resolver.want(row.graph);
588 }
589 }
590 }
591 if let Some(r) = self.resolver.request(&self.caps) {
592 return Ok(Step::Execute(r));
593 }
594 std::mem::take(&mut self.rows)
595 .into_iter()
596 .map(|row| {
597 let name = if row.graph == DEFAULT_GRAPH_ID {
598 GraphName::DefaultGraph
599 } else {
600 crate::encoding::to_graph_name(
601 row.graph,
602 Some(self.resolver.get(row.graph)?),
603 )?
604 };
605 Ok((row, name))
606 })
607 .collect::<Result<Vec<_>>>()
608 .map(Step::Done)
609 }
610 }
611 Registry {
612 request: Some(crate::registry::load_request(caps)),
613 caps: caps.clone(),
614 resolver: TermResolver::default(),
615 rows: Vec::new(),
616 started: false,
617 }
618}
619
620pub fn shape_index_job(caps: &Capabilities) -> OneShot<crate::shapes::ShapeIndex> {
622 OneShot::new(crate::shapes::ShapeIndex::load_request(caps), |r| {
623 crate::shapes::ShapeIndex::from_response(&r)
624 })
625}
626
627fn scalar_changes(r: &Response, n: usize) -> u64 {
629 r.iter().take(n).map(|rs| rs.changes).sum()
630}
631
632pub fn encode_term(t: &Term) -> i64 {
634 term_id(t.as_ref())
635}
636
637pub fn materialize_job(max_rounds: usize, caps: &Capabilities) -> impl Job<Output = u64> {
642 struct Materialize {
643 reset: Option<Request>,
644 round: usize,
645 max_rounds: usize,
646 counting: bool,
647 }
648 impl Job for Materialize {
649 type Output = u64;
650 fn step(&mut self, response: Option<Response>) -> Result<Step<u64>> {
651 if let Some(r) = self.reset.take() {
652 return Ok(Step::Execute(r));
653 }
654 if self.counting {
655 return Ok(Step::Done(
656 scalar(&response.unwrap_or_default()).max(0) as u64
657 ));
658 }
659 let changed: u64 = response.iter().flatten().map(|rs| rs.changes).sum();
660 if self.round > 0 && (changed == 0 || self.round >= self.max_rounds) {
661 self.counting = true;
662 return Ok(Step::Execute(Request::read(vec![
663 "SELECT COUNT(*) FROM quads_inf".into(),
664 ])));
665 }
666 self.round += 1;
667 Ok(Step::Execute(Request::atomic(
668 crate::reason::materialize_round(),
669 )))
670 }
671 }
672 Materialize {
673 reset: Some(Request::atomic(crate::reason::materialize_reset(caps))),
674 round: 0,
675 max_rounds,
676 counting: false,
677 }
678}
679
680pub fn clear_inferences_job() -> OneShot<()> {
682 OneShot::new(
683 Request::atomic(vec![
684 Statement::new("DELETE FROM quads_inf"),
685 Statement::new("DELETE FROM quads_inf_src"),
686 ]),
687 |_| Ok(()),
688 )
689}
690
691pub fn clear_inferences_of_job(producer: &str) -> OneShot<()> {
693 OneShot::new(
694 Request::atomic(crate::reason::inference_reset(producer)),
695 |_| Ok(()),
696 )
697}
698
699pub fn inference_producers_job(quad: QuadRef<'_>) -> OneShot<Vec<String>> {
702 let [s, p, o, g] = [
703 subject_id(quad.subject),
704 named_node_id(quad.predicate.as_str()),
705 term_id(quad.object),
706 graph_id(quad.graph_name),
707 ];
708 OneShot::new(
709 Request::read(vec![Statement::new(format!(
710 "SELECT p.name FROM quads_inf_src q JOIN inf_producers p ON p.id = q.src \
711 WHERE q.s = {s} AND q.p = {p} AND q.o = {o} AND q.g = {g} ORDER BY p.name"
712 ))]),
713 |r| {
714 Ok(r.first()
715 .map(|rs| {
716 rs.rows
717 .iter()
718 .filter_map(|row| match row.first() {
719 Some(SqlValue::Text(t)) => Some(t.clone()),
720 _ => None,
721 })
722 .collect()
723 })
724 .unwrap_or_default())
725 },
726 )
727}