Skip to main content

oxilite_core/
writer.rs

1//! Write path: turns quads into chunked, self-contained SQL statements.
2//!
3//! Term ids are computed in Rust, so an insert never needs to read anything back: terms and
4//! quads go into the database in the same batch, with `INSERT OR IGNORE` making writes
5//! idempotent. Statements are split to respect the backend's maximum SQL length.
6//!
7// @lat: [[architecture#Write path]]
8
9use crate::encoding::{EncodedRows, DEFAULT_GRAPH_ID};
10use crate::sql::{quote_str, sql_f64, Capabilities, Statement};
11use oxrdf::QuadRef;
12use std::fmt::Write;
13
14/// Accumulates `VALUES` tuples and splits them into statements below a size limit.
15struct Chunker<'a> {
16    prefix: &'a str,
17    suffix: &'a str,
18    max: usize,
19    current: String,
20    out: Vec<Statement>,
21}
22
23impl<'a> Chunker<'a> {
24    fn new(prefix: &'a str, suffix: &'a str, max: usize) -> Self {
25        Self {
26            prefix,
27            suffix,
28            max,
29            current: String::new(),
30            out: Vec::new(),
31        }
32    }
33
34    fn push(&mut self, tuple: &str) {
35        if !self.current.is_empty()
36            && self.current.len() + tuple.len() + self.suffix.len() + 1 > self.max
37        {
38            self.flush();
39        }
40        if self.current.is_empty() {
41            self.current.push_str(self.prefix);
42        } else {
43            self.current.push(',');
44        }
45        self.current.push_str(tuple);
46    }
47
48    fn flush(&mut self) {
49        if !self.current.is_empty() {
50            let mut sql = std::mem::take(&mut self.current);
51            sql.push_str(self.suffix);
52            self.out.push(Statement::new(sql));
53        }
54    }
55
56    fn finish(mut self) -> Vec<Statement> {
57        self.flush();
58        self.out
59    }
60}
61
62fn opt_str(out: &mut String, v: Option<&str>) {
63    match v {
64        Some(v) => quote_str(out, v),
65        None => out.push_str("NULL"),
66    }
67}
68
69/// Statements inserting the rows needed to decode encoded terms.
70pub fn term_statements(rows: &EncodedRows, caps: &Capabilities) -> Vec<Statement> {
71    let mut out = Vec::new();
72    let mut terms = Chunker::new(
73        "INSERT OR IGNORE INTO terms(id, lex, dt, lang, dir, num, nt, ts) VALUES ",
74        "",
75        caps.max_sql_len,
76    );
77    let mut tuple = String::new();
78    for r in &rows.terms {
79        tuple.clear();
80        let _ = write!(tuple, "({},", r.id);
81        quote_str(&mut tuple, &r.lex);
82        tuple.push(',');
83        opt_str(&mut tuple, r.dt.as_deref());
84        tuple.push(',');
85        opt_str(&mut tuple, r.lang.as_deref());
86        let _ = write!(
87            tuple,
88            ",{},{},{},{})",
89            r.dir.map_or_else(|| "NULL".into(), |v| v.to_string()),
90            r.num.map_or_else(|| "NULL".into(), sql_f64),
91            r.nt.map_or_else(|| "NULL".into(), |v| v.to_string()),
92            r.ts.map_or_else(|| "NULL".into(), sql_f64),
93        );
94        terms.push(&tuple);
95    }
96    out.extend(terms.finish());
97    let mut triples = Chunker::new(
98        "INSERT OR IGNORE INTO triple_terms(id, s, p, o, vk, sk) VALUES ",
99        "",
100        caps.max_sql_len,
101    );
102    for t in &rows.triples {
103        let mut tuple = format!("({},{},{},{},", t.id, t.s, t.p, t.o);
104        quote_str(&mut tuple, &t.vk);
105        tuple.push(',');
106        quote_str(&mut tuple, &t.sk);
107        tuple.push(')');
108        triples.push(&tuple);
109    }
110    out.extend(triples.finish());
111    out
112}
113
114/// Encoded quads ready to be written.
115#[derive(Debug, Default)]
116pub struct EncodedQuads {
117    pub rows: EncodedRows,
118    pub quads: Vec<[i64; 4]>,
119}
120
121impl EncodedQuads {
122    pub fn new<'a>(quads: impl IntoIterator<Item = QuadRef<'a>>) -> Self {
123        let mut me = Self::default();
124        for q in quads {
125            let ids = me.rows.quad(q);
126            me.quads.push(ids);
127        }
128        me.rows.dedup();
129        me
130    }
131
132    /// Statements inserting the quads (and their terms / graph names).
133    pub fn insert_statements(&self, caps: &Capabilities) -> Vec<Statement> {
134        let mut out = term_statements(&self.rows, caps);
135        let mut graphs: Vec<i64> = self
136            .quads
137            .iter()
138            .map(|q| q[3])
139            .filter(|g| *g != DEFAULT_GRAPH_ID)
140            .collect();
141        graphs.sort_unstable();
142        graphs.dedup();
143        let mut g = Chunker::new(
144            "INSERT OR IGNORE INTO graphs(id) VALUES ",
145            "",
146            caps.max_sql_len,
147        );
148        for id in graphs {
149            g.push(&format!("({id})"));
150        }
151        out.extend(g.finish());
152        out.extend(quad_insert_statements(&self.quads, caps));
153        out
154    }
155
156    /// Statements deleting the quads.
157    pub fn delete_statements(&self, caps: &Capabilities) -> Vec<Statement> {
158        quad_delete_statements(&self.quads, caps)
159    }
160}
161
162/// One atomic request: `prefix`, the inserts of `quads`, then `suffix`. Atomicity forbids
163/// splitting it, so `Err(n)` reports its `n` statements when they exceed the backend's
164/// per-request limit (each statement is already below the SQL-length limit).
165pub fn atomic_request(
166    prefix: Vec<Statement>,
167    quads: &EncodedQuads,
168    suffix: Vec<Statement>,
169    caps: &Capabilities,
170) -> std::result::Result<crate::sql::Request, usize> {
171    let mut s = prefix;
172    s.extend(quads.insert_statements(caps));
173    s.extend(suffix);
174    if s.len() > caps.max_statements {
175        return Err(s.len());
176    }
177    Ok(crate::sql::Request::atomic(s))
178}
179
180/// `INSERT OR IGNORE INTO quads` statements.
181pub fn quad_insert_statements(quads: &[[i64; 4]], caps: &Capabilities) -> Vec<Statement> {
182    insert_statements_into("quads", quads, caps)
183}
184
185/// `INSERT OR IGNORE INTO <table>(s, p, o, g)` statements (`quads` or `quads_inf`).
186pub fn insert_statements_into(
187    table: &str,
188    quads: &[[i64; 4]],
189    caps: &Capabilities,
190) -> Vec<Statement> {
191    let prefix = format!("INSERT OR IGNORE INTO {table}(s, p, o, g) VALUES ");
192    let mut c = Chunker::new(&prefix, "", caps.max_sql_len);
193    for [s, p, o, g] in quads {
194        c.push(&format!("({s},{p},{o},{g})"));
195    }
196    c.finish()
197}
198
199/// `DELETE FROM quads` statements.
200pub fn quad_delete_statements(quads: &[[i64; 4]], caps: &Capabilities) -> Vec<Statement> {
201    if let [[s, p, o, g]] = quads {
202        return vec![Statement::new(format!(
203            "DELETE FROM quads WHERE s = {s} AND p = {p} AND o = {o} AND g = {g}"
204        ))];
205    }
206    let mut c = Chunker::new(
207        "DELETE FROM quads WHERE (s, p, o, g) IN (VALUES ",
208        ")",
209        caps.max_sql_len,
210    );
211    for [s, p, o, g] in quads {
212        c.push(&format!("({s},{p},{o},{g})"));
213    }
214    c.finish()
215}
216
217#[cfg(test)]
218mod tests {
219    use super::*;
220    use oxrdf::{GraphName, Literal, NamedNode, Quad};
221
222    // @lat: [[tests#Write path#Statements respect the size limit]]
223    #[test]
224    fn statements_respect_the_size_limit() {
225        let quads: Vec<Quad> = (0..500)
226            .map(|i| {
227                Quad::new(
228                    NamedNode::new_unchecked(format!("http://example.com/s{i}")),
229                    NamedNode::new_unchecked("http://example.com/p"),
230                    Literal::new_simple_literal(format!("value number {i} with 'quotes'")),
231                    GraphName::DefaultGraph,
232                )
233            })
234            .collect();
235        let enc = EncodedQuads::new(quads.iter().map(Quad::as_ref));
236        let caps = Capabilities {
237            max_sql_len: 2_000,
238            ..Capabilities::d1()
239        };
240        let stmts = enc.insert_statements(&caps);
241        assert!(stmts.len() > 5);
242        assert!(stmts.iter().all(|s| s.sql.len() <= 2_000));
243        assert!(stmts.iter().all(|s| s.params.is_empty()));
244    }
245}