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
85fn count_tail(r: &Response, n: usize) -> u64 {
86    r.iter().rev().take(n).map(|rs| rs.changes).sum()
87}
88
89/// Atomically inserts quads; returns how many were new.
90pub fn insert_job<'a>(
91    quads: impl IntoIterator<Item = QuadRef<'a>>,
92    caps: &Capabilities,
93) -> OneShot<u64> {
94    let quads: Vec<QuadRef<'a>> = quads.into_iter().collect();
95    let schema = quads.iter().any(|q| crate::reason::is_schema_quad(*q));
96    let enc = EncodedQuads::new(quads);
97    let quad_stmts = crate::writer::quad_insert_statements(&enc.quads, caps).len();
98    let mut stmts = enc.insert_statements(caps);
99    let end = stmts.len();
100    if schema {
101        stmts.extend(crate::reason::closure_statements());
102    }
103    OneShot::new(Request::atomic(stmts), move |mut r| {
104        r.truncate(end);
105        Ok(count_tail(&r, quad_stmts))
106    })
107}
108
109/// Insert statements for a chunk (bulk loading).
110pub fn insert_request<'a>(
111    quads: impl IntoIterator<Item = QuadRef<'a>>,
112    caps: &Capabilities,
113) -> Request {
114    Request::atomic(EncodedQuads::new(quads).insert_statements(caps))
115}
116
117/// Atomically removes quads; returns how many were present.
118pub fn remove_job<'a>(
119    quads: impl IntoIterator<Item = QuadRef<'a>>,
120    caps: &Capabilities,
121) -> OneShot<u64> {
122    let quads: Vec<QuadRef<'a>> = quads.into_iter().collect();
123    let schema = quads.iter().any(|q| crate::reason::is_schema_quad(*q));
124    let enc = EncodedQuads::new(quads);
125    let mut stmts = enc.delete_statements(caps);
126    let end = stmts.len();
127    if schema {
128        stmts.extend(crate::reason::closure_statements());
129    }
130    OneShot::new(Request::atomic(stmts), move |r| {
131        Ok(r.iter().take(end).map(|rs| rs.changes).sum())
132    })
133}
134
135fn scalar(r: &Response) -> i64 {
136    r.first()
137        .and_then(|rs| rs.rows.first())
138        .and_then(|row| row.first())
139        .and_then(SqlValue::as_i64)
140        .unwrap_or(0)
141}
142
143pub fn contains_job(quad: QuadRef<'_>) -> OneShot<bool> {
144    let [s, p, o, g] = [
145        subject_id(quad.subject),
146        named_node_id(quad.predicate.as_str()),
147        term_id(quad.object),
148        graph_id(quad.graph_name),
149    ];
150    OneShot::new(
151        Request::read(vec![Statement::new(format!(
152            "SELECT EXISTS (SELECT 1 FROM quads WHERE s = {s} AND p = {p} AND o = {o} AND g = {g})"
153        ))]),
154        |r| Ok(scalar(&r) != 0),
155    )
156}
157
158pub fn len_job() -> OneShot<usize> {
159    OneShot::new(
160        Request::read(vec!["SELECT COUNT(*) FROM quads".into()]),
161        |r| Ok(scalar(&r) as usize),
162    )
163}
164
165pub fn is_empty_job() -> OneShot<bool> {
166    OneShot::new(
167        Request::read(vec!["SELECT NOT EXISTS (SELECT 1 FROM quads)".into()]),
168        |r| Ok(scalar(&r) != 0),
169    )
170}
171
172/// A job returning quads matching a SQL WHERE clause, with batched term resolution.
173pub struct ScanJob {
174    sql: Option<String>,
175    caps: Capabilities,
176    resolver: TermResolver,
177    rows: Vec<[i64; 4]>,
178    started: bool,
179}
180
181impl ScanJob {
182    fn new(where_clause: String, caps: &Capabilities) -> Self {
183        let sql = format!(
184            "SELECT {}, {}, {}, {} FROM quads{}",
185            id_col(caps, "s"),
186            id_col(caps, "p"),
187            id_col(caps, "o"),
188            id_col(caps, "g"),
189            if where_clause.is_empty() {
190                String::new()
191            } else {
192                format!(" WHERE {where_clause}")
193            }
194        );
195        Self {
196            sql: Some(sql),
197            caps: caps.clone(),
198            resolver: TermResolver::default(),
199            rows: Vec::new(),
200            started: false,
201        }
202    }
203
204    fn finish(&self) -> Result<Vec<Quad>> {
205        self.rows
206            .iter()
207            .map(|[s, p, o, g]| {
208                let gname = if *g == DEFAULT_GRAPH_ID {
209                    GraphName::DefaultGraph
210                } else {
211                    crate::encoding::to_graph_name(*g, Some(self.resolver.get(*g)?))?
212                };
213                crate::encoding::make_quad(
214                    self.resolver.get(*s)?,
215                    self.resolver.get(*p)?,
216                    self.resolver.get(*o)?,
217                    gname,
218                )
219            })
220            .collect()
221    }
222}
223
224impl Job for ScanJob {
225    type Output = Vec<Quad>;
226
227    fn step(&mut self, response: Option<Response>) -> Result<Step<Vec<Quad>>> {
228        if let Some(sql) = self.sql.take() {
229            return Ok(Step::Execute(Request::read(vec![Statement::new(sql)])));
230        }
231        let response =
232            response.ok_or_else(|| Error::Other("scan resumed without response".into()))?;
233        if self.started {
234            self.resolver.absorb(response)?;
235        } else {
236            self.started = true;
237            for rs in response {
238                for row in rs.rows {
239                    let ids: Vec<i64> = row.iter().filter_map(SqlValue::as_i64).collect();
240                    let [s, p, o, g] = ids[..] else {
241                        return Err(Error::corrupted("bad quad row"));
242                    };
243                    for id in [s, p, o, g] {
244                        self.resolver.want(id);
245                    }
246                    self.rows.push([s, p, o, g]);
247                }
248            }
249        }
250        match self.resolver.request(&self.caps) {
251            Some(r) => Ok(Step::Execute(r)),
252            None => self.finish().map(Step::Done),
253        }
254    }
255}
256
257/// Quads whose subject (or, with `incoming`, object) is one of `nodes`, optionally restricted
258/// to some predicates and to the default graph: one neighbourhood hop for bounded prefetches.
259pub fn neighbourhood_job(
260    nodes: &[TermRef<'_>],
261    incoming: bool,
262    predicates: Option<&[NamedNodeRef<'_>]>,
263    default_graph_only: bool,
264    caps: &Capabilities,
265) -> ScanJob {
266    let ids: Vec<String> = nodes.iter().map(|t| term_id(*t).to_string()).collect();
267    let mut w = vec![format!(
268        "{} IN ({})",
269        if incoming { "o" } else { "s" },
270        if ids.is_empty() {
271            "NULL".into()
272        } else {
273            ids.join(",")
274        }
275    )];
276    if let Some(ps) = predicates {
277        let ps: Vec<String> = ps
278            .iter()
279            .map(|p| named_node_id(p.as_str()).to_string())
280            .collect();
281        w.push(format!(
282            "p IN ({})",
283            if ps.is_empty() {
284                "NULL".into()
285            } else {
286                ps.join(",")
287            }
288        ));
289    }
290    if default_graph_only {
291        w.push(format!("g = {DEFAULT_GRAPH_ID}"));
292    }
293    ScanJob::new(w.join(" AND "), caps)
294}
295
296/// `quads_for_pattern` as a job. `graph_name = None` matches every graph, default included.
297pub fn scan_job(
298    subject: Option<NamedOrBlankNodeRef<'_>>,
299    predicate: Option<NamedNodeRef<'_>>,
300    object: Option<TermRef<'_>>,
301    graph_name: Option<GraphNameRef<'_>>,
302    caps: &Capabilities,
303) -> ScanJob {
304    let mut w = Vec::new();
305    if let Some(s) = subject {
306        w.push(format!("s = {}", subject_id(s)));
307    }
308    if let Some(p) = predicate {
309        w.push(format!("p = {}", named_node_id(p.as_str())));
310    }
311    if let Some(o) = object {
312        w.push(format!("o = {}", term_id(o)));
313    }
314    if let Some(g) = graph_name {
315        w.push(format!("g = {}", graph_id(g)));
316    }
317    ScanJob::new(w.join(" AND "), caps)
318}
319
320/// Lists named graphs.
321pub fn named_graphs_job(caps: &Capabilities) -> impl Job<Output = Vec<NamedOrBlankNode>> {
322    struct Graphs {
323        sql: Option<String>,
324        caps: Capabilities,
325        resolver: TermResolver,
326        ids: Vec<i64>,
327        started: bool,
328    }
329    impl Job for Graphs {
330        type Output = Vec<NamedOrBlankNode>;
331        fn step(&mut self, response: Option<Response>) -> Result<Step<Self::Output>> {
332            if let Some(sql) = self.sql.take() {
333                return Ok(Step::Execute(Request::read(vec![Statement::new(sql)])));
334            }
335            let response = response.unwrap_or_default();
336            if self.started {
337                self.resolver.absorb(response)?;
338            } else {
339                self.started = true;
340                self.ids = crate::resolve::ids_of(&response, 0);
341                for id in &self.ids {
342                    self.resolver.want(*id);
343                }
344            }
345            if let Some(r) = self.resolver.request(&self.caps) {
346                return Ok(Step::Execute(r));
347            }
348            self.ids
349                .iter()
350                .map(|id| crate::encoding::to_subject(self.resolver.get(*id)?))
351                .collect::<Result<_>>()
352                .map(Step::Done)
353        }
354    }
355    Graphs {
356        sql: Some(format!("SELECT {} FROM graphs", id_col(caps, "id"))),
357        caps: caps.clone(),
358        resolver: TermResolver::default(),
359        ids: Vec::new(),
360        started: false,
361    }
362}
363
364pub fn contains_named_graph_job(g: NamedOrBlankNodeRef<'_>) -> OneShot<bool> {
365    let id = subject_id(g);
366    OneShot::new(
367        Request::read(vec![Statement::new(format!(
368            "SELECT EXISTS (SELECT 1 FROM graphs WHERE id = {id})"
369        ))]),
370        |r| Ok(scalar(&r) != 0),
371    )
372}
373
374/// Adds a named graph; returns whether it was new.
375pub fn insert_named_graph_job(g: NamedOrBlankNodeRef<'_>, caps: &Capabilities) -> OneShot<bool> {
376    let mut rows = EncodedRows::default();
377    let id = rows.subject(g);
378    let mut stmts = term_statements(&rows, caps);
379    stmts.push(Statement::new(format!(
380        "INSERT OR IGNORE INTO graphs(id) VALUES ({id})"
381    )));
382    OneShot::new(Request::atomic(stmts), |r| Ok(count_tail(&r, 1) > 0))
383}
384
385/// Removes a named graph and its quads; returns whether it existed.
386pub fn remove_named_graph_job(g: NamedOrBlankNodeRef<'_>) -> OneShot<bool> {
387    let id = subject_id(g);
388    let mut stmts = vec![
389        Statement::new(format!("DELETE FROM quads WHERE g = {id}")),
390        Statement::new(format!("DELETE FROM graphs WHERE id = {id}")),
391    ];
392    stmts.extend(crate::reason::closure_statements());
393    OneShot::new(Request::atomic(stmts), |r| {
394        Ok(r.iter().take(2).any(|rs| rs.changes > 0))
395    })
396}
397
398/// Removes all quads of a graph (keeps the graph name).
399pub fn clear_graph_job(g: GraphNameRef<'_>) -> OneShot<()> {
400    let id = graph_id(g);
401    let mut stmts = vec![Statement::new(format!("DELETE FROM quads WHERE g = {id}"))];
402    stmts.extend(crate::reason::closure_statements());
403    OneShot::new(Request::atomic(stmts), |_| Ok(()))
404}
405
406/// Removes everything.
407pub fn clear_job() -> OneShot<()> {
408    OneShot::new(
409        Request::atomic(vec![
410            "DELETE FROM quads".into(),
411            "DELETE FROM quads_inf".into(),
412            "DELETE FROM tbox_closure".into(),
413            "DELETE FROM graphs".into(),
414            "DELETE FROM triple_terms".into(),
415            "DELETE FROM terms".into(),
416        ]),
417        |_| Ok(()),
418    )
419}
420
421/// Helper used by drivers: encodes a term to its id without I/O.
422pub fn encode_term(t: &Term) -> i64 {
423    term_id(t.as_ref())
424}
425
426/// Recomputes the OWL 2 RL materialization (`quads_inf`): discards previous inferences, then
427/// runs rule rounds (one atomic request each, i.e. one D1 batch) until a round infers nothing.
428/// Returns the number of inferred triples. Rounds are not one transaction: readers may see a
429/// partial materialization while it runs.
430pub fn materialize_job(max_rounds: usize, caps: &Capabilities) -> impl Job<Output = u64> {
431    struct Materialize {
432        reset: Option<Request>,
433        round: usize,
434        max_rounds: usize,
435        counting: bool,
436    }
437    impl Job for Materialize {
438        type Output = u64;
439        fn step(&mut self, response: Option<Response>) -> Result<Step<u64>> {
440            if let Some(r) = self.reset.take() {
441                return Ok(Step::Execute(r));
442            }
443            if self.counting {
444                return Ok(Step::Done(
445                    scalar(&response.unwrap_or_default()).max(0) as u64
446                ));
447            }
448            let changed: u64 = response.iter().flatten().map(|rs| rs.changes).sum();
449            if self.round > 0 && (changed == 0 || self.round >= self.max_rounds) {
450                self.counting = true;
451                return Ok(Step::Execute(Request::read(vec![
452                    "SELECT COUNT(*) FROM quads_inf".into(),
453                ])));
454            }
455            self.round += 1;
456            Ok(Step::Execute(Request::atomic(
457                crate::reason::materialize_round(),
458            )))
459        }
460    }
461    Materialize {
462        reset: Some(Request::atomic(crate::reason::materialize_reset(caps))),
463        round: 0,
464        max_rounds,
465        counting: false,
466    }
467}
468
469/// Removes every materialized inference.
470pub fn clear_inferences_job() -> OneShot<()> {
471    OneShot::new(
472        Request::atomic(vec![Statement::new("DELETE FROM quads_inf")]),
473        |_| Ok(()),
474    )
475}