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