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::{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
28/// Creates the schema, then loads statistics.
29///
30/// The versioning level recorded in the store wins: opening never changes it, except that an
31/// empty store opened with a higher level gets it (that is how a store is created versioned).
32/// A higher level on a store holding data is refused: raising it is an explicit level change.
33pub 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
110/// Loads statistics only.
111pub fn stats_job(caps: &Capabilities) -> OneShot<Stats> {
112    OneShot::new(Stats::load_request(caps), |r| Stats::from_response(&r))
113}
114
115/// Recomputes and reloads statistics.
116pub 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
141/// The statements that rebuild the derived schema caches inside a write's own request:
142/// `tbox_closure` when a schema axiom changed, the shape index when a SHACL triple did.
143///
144/// Every driver that assembles a write request itself (the async store, the JavaScript
145/// bindings, the Cypher writer) appends this, so the caches cannot be refreshed on one backend
146/// and forgotten on another.
147pub 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
175/// Atomically inserts quads; returns how many were new.
176pub 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
193/// Insert statements for a chunk (bulk loading).
194pub 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
201/// Atomically removes quads; returns how many were present.
202pub 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
254/// A job returning quads matching a SQL WHERE clause, with batched term resolution.
255pub 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
339/// Quads whose subject (or, with `incoming`, object) is one of `nodes`, optionally restricted
340/// to some predicates and to the default graph: one neighbourhood hop for bounded prefetches.
341pub 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
378/// `quads_for_pattern` as a job. `graph_name = None` matches every graph, default included.
379pub 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
402/// Lists named graphs.
403pub 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
456/// Adds a named graph; returns whether it was new.
457pub 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
467/// Removes a named graph and its quads; returns whether it existed.
468pub 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
481/// Removes all quads of a graph (keeps the graph name).
482pub 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
489/// Removes everything.
490pub 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
509/// Registers a schema graph, then rebuilds what its role feeds.
510///
511/// A registration changes the *scope* of the closure and the shape index — not just their
512/// content — so both are rebuilt in the same atomic request as the registration itself.
513///
514/// A named graph may be registered before it holds anything, so the registration creates it
515/// (as `insert_named_graph` would): otherwise the registry would name a graph the term
516/// dictionary has never heard of.
517pub 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
541/// Removes a registration, keeping the graph's triples.
542pub 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
548/// Activates or deactivates a registration.
549pub 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
555/// Removes a registration together with every quad of its graph.
556pub 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
562/// Reads the registry, with each row's graph name resolved.
563pub 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
620/// Reads the compiled shape index.
621pub 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
627/// Rows changed by the first `n` statements of a response.
628fn scalar_changes(r: &Response, n: usize) -> u64 {
629    r.iter().take(n).map(|rs| rs.changes).sum()
630}
631
632/// Helper used by drivers: encodes a term to its id without I/O.
633pub fn encode_term(t: &Term) -> i64 {
634    term_id(t.as_ref())
635}
636
637/// Recomputes the OWL 2 RL materialization (`quads_inf`): discards previous inferences, then
638/// runs rule rounds (one atomic request each, i.e. one D1 batch) until a round infers nothing.
639/// Returns the number of inferred triples. Rounds are not one transaction: readers may see a
640/// partial materialization while it runs.
641pub 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
680/// Removes every materialized inference.
681pub 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
691/// Removes the inferences of one producer, keeping what other producers also derived.
692pub fn clear_inferences_of_job(producer: &str) -> OneShot<()> {
693    OneShot::new(
694        Request::atomic(crate::reason::inference_reset(producer)),
695        |_| Ok(()),
696    )
697}
698
699/// The names of the producers that derived a quad (empty when it is not inferred). The
700/// default graph matches the graph-0 conclusions every producer writes.
701pub 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}