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 tbox_closure".into(),
440            "DELETE FROM shapes_index".into(),
441            "DELETE FROM shapes_in".into(),
442            "DELETE FROM schema_graphs".into(),
443            "DELETE FROM graphs".into(),
444            "DELETE FROM triple_terms".into(),
445            "DELETE FROM terms".into(),
446        ]),
447        |_| Ok(()),
448    )
449}
450
451/// Registers a schema graph, then rebuilds what its role feeds.
452///
453/// A registration changes the *scope* of the closure and the shape index — not just their
454/// content — so both are rebuilt in the same atomic request as the registration itself.
455///
456/// A named graph may be registered before it holds anything, so the registration creates it
457/// (as `insert_named_graph` would): otherwise the registry would name a graph the term
458/// dictionary has never heard of.
459pub fn register_schema_graph_job(
460    entry: &crate::registry::SchemaGraph,
461    graph: GraphNameRef<'_>,
462    caps: &Capabilities,
463) -> OneShot<()> {
464    let mut stmts = Vec::new();
465    let named: Option<NamedOrBlankNodeRef<'_>> = match graph {
466        GraphNameRef::NamedNode(n) => Some(n.into()),
467        GraphNameRef::BlankNode(b) => Some(b.into()),
468        GraphNameRef::DefaultGraph => None,
469    };
470    if let Some(g) = named {
471        let mut rows = EncodedRows::default();
472        let id = rows.subject(g);
473        stmts.extend(term_statements(&rows, caps));
474        stmts.push(Statement::new(format!(
475            "INSERT OR IGNORE INTO graphs(id) VALUES ({id})"
476        )));
477    }
478    stmts.extend(crate::registry::register_statements(entry));
479    stmts.extend(schema_refresh(true, true));
480    OneShot::new(Request::atomic(stmts), |_| Ok(()))
481}
482
483/// Removes a registration, keeping the graph's triples.
484pub fn unregister_schema_graph_job(graph: i64) -> OneShot<bool> {
485    let mut stmts = crate::registry::unregister_statements(graph);
486    stmts.extend(schema_refresh(true, true));
487    OneShot::new(Request::atomic(stmts), |r| Ok(scalar_changes(&r, 1) > 0))
488}
489
490/// Activates or deactivates a registration.
491pub fn set_schema_graph_active_job(graph: i64, active: bool) -> OneShot<bool> {
492    let mut stmts = crate::registry::set_active_statements(graph, active);
493    stmts.extend(schema_refresh(true, true));
494    OneShot::new(Request::atomic(stmts), |r| Ok(scalar_changes(&r, 1) > 0))
495}
496
497/// Removes a registration together with every quad of its graph.
498pub fn drop_schema_graph_job(graph: i64) -> OneShot<u64> {
499    let mut stmts = crate::registry::drop_statements(graph);
500    stmts.extend(schema_refresh(true, true));
501    OneShot::new(Request::atomic(stmts), |r| Ok(scalar_changes(&r, 1)))
502}
503
504/// Reads the registry, with each row's graph name resolved.
505pub fn schema_graphs_job(
506    caps: &Capabilities,
507) -> impl Job<Output = Vec<(crate::registry::SchemaGraph, GraphName)>> {
508    struct Registry {
509        request: Option<Request>,
510        caps: Capabilities,
511        resolver: TermResolver,
512        rows: Vec<crate::registry::SchemaGraph>,
513        started: bool,
514    }
515    impl Job for Registry {
516        type Output = Vec<(crate::registry::SchemaGraph, GraphName)>;
517        fn step(&mut self, response: Option<Response>) -> Result<Step<Self::Output>> {
518            if let Some(r) = self.request.take() {
519                return Ok(Step::Execute(r));
520            }
521            let response = response.unwrap_or_default();
522            if self.started {
523                self.resolver.absorb(response)?;
524            } else {
525                self.started = true;
526                self.rows = crate::registry::from_response(&response)?;
527                for row in &self.rows {
528                    if row.graph != DEFAULT_GRAPH_ID {
529                        self.resolver.want(row.graph);
530                    }
531                }
532            }
533            if let Some(r) = self.resolver.request(&self.caps) {
534                return Ok(Step::Execute(r));
535            }
536            std::mem::take(&mut self.rows)
537                .into_iter()
538                .map(|row| {
539                    let name = if row.graph == DEFAULT_GRAPH_ID {
540                        GraphName::DefaultGraph
541                    } else {
542                        crate::encoding::to_graph_name(
543                            row.graph,
544                            Some(self.resolver.get(row.graph)?),
545                        )?
546                    };
547                    Ok((row, name))
548                })
549                .collect::<Result<Vec<_>>>()
550                .map(Step::Done)
551        }
552    }
553    Registry {
554        request: Some(crate::registry::load_request(caps)),
555        caps: caps.clone(),
556        resolver: TermResolver::default(),
557        rows: Vec::new(),
558        started: false,
559    }
560}
561
562/// Reads the compiled shape index.
563pub fn shape_index_job(caps: &Capabilities) -> OneShot<crate::shapes::ShapeIndex> {
564    OneShot::new(crate::shapes::ShapeIndex::load_request(caps), |r| {
565        crate::shapes::ShapeIndex::from_response(&r)
566    })
567}
568
569/// Rows changed by the first `n` statements of a response.
570fn scalar_changes(r: &Response, n: usize) -> u64 {
571    r.iter().take(n).map(|rs| rs.changes).sum()
572}
573
574/// Helper used by drivers: encodes a term to its id without I/O.
575pub fn encode_term(t: &Term) -> i64 {
576    term_id(t.as_ref())
577}
578
579/// Recomputes the OWL 2 RL materialization (`quads_inf`): discards previous inferences, then
580/// runs rule rounds (one atomic request each, i.e. one D1 batch) until a round infers nothing.
581/// Returns the number of inferred triples. Rounds are not one transaction: readers may see a
582/// partial materialization while it runs.
583pub fn materialize_job(max_rounds: usize, caps: &Capabilities) -> impl Job<Output = u64> {
584    struct Materialize {
585        reset: Option<Request>,
586        round: usize,
587        max_rounds: usize,
588        counting: bool,
589    }
590    impl Job for Materialize {
591        type Output = u64;
592        fn step(&mut self, response: Option<Response>) -> Result<Step<u64>> {
593            if let Some(r) = self.reset.take() {
594                return Ok(Step::Execute(r));
595            }
596            if self.counting {
597                return Ok(Step::Done(
598                    scalar(&response.unwrap_or_default()).max(0) as u64
599                ));
600            }
601            let changed: u64 = response.iter().flatten().map(|rs| rs.changes).sum();
602            if self.round > 0 && (changed == 0 || self.round >= self.max_rounds) {
603                self.counting = true;
604                return Ok(Step::Execute(Request::read(vec![
605                    "SELECT COUNT(*) FROM quads_inf".into(),
606                ])));
607            }
608            self.round += 1;
609            Ok(Step::Execute(Request::atomic(
610                crate::reason::materialize_round(),
611            )))
612        }
613    }
614    Materialize {
615        reset: Some(Request::atomic(crate::reason::materialize_reset(caps))),
616        round: 0,
617        max_rounds,
618        counting: false,
619    }
620}
621
622/// Removes every materialized inference.
623pub fn clear_inferences_job() -> OneShot<()> {
624    OneShot::new(
625        Request::atomic(vec![Statement::new("DELETE FROM quads_inf")]),
626        |_| Ok(()),
627    )
628}