Skip to main content

oxilite_core/
ops.rs

1//! Store operations as sans-IO jobs, shared by the sync, async and JavaScript drivers.
2//!
3// @lat: [[architecture#Sans-IO core]]
4
5use 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
28/// Creates the schema, then loads statistics.
29pub 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
54/// Loads statistics only.
55pub fn stats_job(caps: &Capabilities) -> OneShot<Stats> {
56    OneShot::new(Stats::load_request(caps), |r| Stats::from_response(&r))
57}
58
59/// Recomputes and reloads statistics.
60pub 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
85/// The statements that rebuild the derived schema caches inside a write's own request:
86/// `tbox_closure` when a schema axiom changed, the shape index when a SHACL triple did.
87///
88/// Every driver that assembles a write request itself (the async store, the JavaScript
89/// bindings, the Cypher writer) appends this, so the caches cannot be refreshed on one backend
90/// and forgotten on another.
91pub 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
119/// Atomically inserts quads; returns how many were new.
120pub 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
137/// Insert statements for a chunk (bulk loading).
138pub 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
145/// Atomically removes quads; returns how many were present.
146pub 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
198/// A job returning quads matching a SQL WHERE clause, with batched term resolution.
199pub 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
283/// Quads whose subject (or, with `incoming`, object) is one of `nodes`, optionally restricted
284/// to some predicates and to the default graph: one neighbourhood hop for bounded prefetches.
285pub 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
322/// `quads_for_pattern` as a job. `graph_name = None` matches every graph, default included.
323pub 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
346/// Lists named graphs.
347pub 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
400/// Adds a named graph; returns whether it was new.
401pub 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
411/// Removes a named graph and its quads; returns whether it existed.
412pub 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
425/// Removes all quads of a graph (keeps the graph name).
426pub 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
433/// Removes everything.
434pub 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
453/// Registers a schema graph, then rebuilds what its role feeds.
454///
455/// A registration changes the *scope* of the closure and the shape index — not just their
456/// content — so both are rebuilt in the same atomic request as the registration itself.
457///
458/// A named graph may be registered before it holds anything, so the registration creates it
459/// (as `insert_named_graph` would): otherwise the registry would name a graph the term
460/// dictionary has never heard of.
461pub 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
485/// Removes a registration, keeping the graph's triples.
486pub 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
492/// Activates or deactivates a registration.
493pub 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
499/// Removes a registration together with every quad of its graph.
500pub 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
506/// Reads the registry, with each row's graph name resolved.
507pub 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
564/// Reads the compiled shape index.
565pub 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
571/// Rows changed by the first `n` statements of a response.
572fn scalar_changes(r: &Response, n: usize) -> u64 {
573    r.iter().take(n).map(|rs| rs.changes).sum()
574}
575
576/// Helper used by drivers: encodes a term to its id without I/O.
577pub fn encode_term(t: &Term) -> i64 {
578    term_id(t.as_ref())
579}
580
581/// Recomputes the OWL 2 RL materialization (`quads_inf`): discards previous inferences, then
582/// runs rule rounds (one atomic request each, i.e. one D1 batch) until a round infers nothing.
583/// Returns the number of inferred triples. Rounds are not one transaction: readers may see a
584/// partial materialization while it runs.
585pub 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
624/// Removes every materialized inference.
625pub 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
635/// Removes the inferences of one producer, keeping what other producers also derived.
636pub fn clear_inferences_of_job(producer: &str) -> OneShot<()> {
637    OneShot::new(
638        Request::atomic(crate::reason::inference_reset(producer)),
639        |_| Ok(()),
640    )
641}
642
643/// The names of the producers that derived a quad (empty when it is not inferred). The
644/// default graph matches the graph-0 conclusions every producer writes.
645pub 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}