Skip to main content

oxilite_core/
version.rs

1//! Optional versioning: a store clock, an immutable change log, and time travel.
2//!
3//! A store chooses a [`Versioning`] level. `Off` (the default) is the plain store and costs
4//! nothing. `Stamped` adds a monotonic store clock: every atomic write opens one *tick* (a row of
5//! `ticks`, with its wall time, author and message), and every quad records the tick that added
6//! it in `quads.t`. `Log` adds the immutable change log `quad_log`, written by triggers on
7//! `quads`, from which any past state is reconstructed (`as_of`).
8//!
9//! The tick is opened by the store, not by the writers: [`prepare`] prepends one statement to
10//! every atomic request that writes `quads`, and [`VersionedBackend`] / [`Versioned`] apply it to
11//! every request of a store, so no writer can forget it.
12//!
13// @lat: [[architecture#Versioning]]
14
15use crate::error::{Error, Result};
16use crate::job::{AsyncBackend, Job, Step, SyncBackend};
17use crate::sql::{
18    col, expect_len, sql_opt_str, sql_str, Capabilities, Mode, Request, Response, Statement,
19};
20use std::fmt;
21use std::str::FromStr;
22use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
23use std::sync::RwLock;
24
25/// How much history a store keeps. Levels nest: each keeps everything the lower ones keep.
26#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Default)]
27#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
28#[cfg_attr(feature = "serde", serde(rename_all = "lowercase"))]
29pub enum Versioning {
30    /// No history (the default): the plain store, with no extra table, column or statement.
31    #[default]
32    Off,
33    /// A monotonic store clock: one tick per atomic write, and the tick that added each quad.
34    Stamped,
35    /// An immutable change log: every effective change, and any past state on request.
36    Log,
37}
38
39impl Versioning {
40    pub const ALL: [Versioning; 3] = [Versioning::Off, Versioning::Stamped, Versioning::Log];
41
42    pub fn as_str(self) -> &'static str {
43        match self {
44            Self::Off => "off",
45            Self::Stamped => "stamped",
46            Self::Log => "log",
47        }
48    }
49
50    fn index(self) -> u8 {
51        self as u8
52    }
53
54    fn from_index(i: u8) -> Self {
55        Self::ALL[usize::from(i).min(2)]
56    }
57}
58
59impl fmt::Display for Versioning {
60    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
61        f.write_str(self.as_str())
62    }
63}
64
65impl FromStr for Versioning {
66    type Err = Error;
67    fn from_str(s: &str) -> Result<Self> {
68        match s.to_ascii_lowercase().as_str() {
69            "off" | "none" => Ok(Self::Off),
70            "stamped" => Ok(Self::Stamped),
71            "log" => Ok(Self::Log),
72            _ => Err(Error::Other(format!(
73                "unknown versioning level `{s}` (expected off, stamped or log)"
74            ))),
75        }
76    }
77}
78
79/// Whether a change log exists, and whether it is still recording.
80#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
81#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
82#[cfg_attr(feature = "serde", serde(rename_all = "lowercase"))]
83pub enum History {
84    /// No change log.
85    #[default]
86    None,
87    /// The log records every change (level `log`).
88    Live,
89    /// The log stopped recording at a freeze (the level was lowered); it still answers as-of
90    /// queries up to the freeze.
91    Frozen,
92}
93
94impl History {
95    pub fn as_str(self) -> &'static str {
96        match self {
97            Self::None => "none",
98            Self::Live => "live",
99            Self::Frozen => "frozen",
100        }
101    }
102}
103
104/// The versioning state recorded in `oxilite_meta` (read with the statistics at open).
105#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
106#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
107#[cfg_attr(feature = "serde", serde(default, rename_all = "camelCase"))]
108pub struct VersionState {
109    pub level: Versioning,
110    pub history: History,
111    /// `quads.t` exists (it survives a soft downgrade to `off`).
112    pub stamp_column: bool,
113    /// The optional index on `quads.t` (fast "added since").
114    pub stamp_index: bool,
115    /// The optional `(p,o,s,g,tx)` and `(o,s,p,g,tx)` log indexes (fast as-of patterns).
116    pub as_of_index: bool,
117}
118
119impl VersionState {
120    /// Absorbs one `oxilite_meta` entry.
121    pub fn absorb(&mut self, key: &str, value: &str) {
122        match key {
123            "versioning" => self.level = value.parse().unwrap_or_default(),
124            "history" => {
125                self.history = match value {
126                    "live" => History::Live,
127                    "frozen" => History::Frozen,
128                    _ => History::None,
129                }
130            }
131            "stamp_column" => self.stamp_column = value == "1",
132            "stamp_index" => self.stamp_index = value == "1",
133            "as_of_index" => self.as_of_index = value == "1",
134            _ => {}
135        }
136    }
137
138    /// Is there a clock (`ticks`) to read?
139    pub fn has_ticks(&self) -> bool {
140        self.level >= Versioning::Stamped || self.stamp_column || self.history != History::None
141    }
142}
143
144/// Who made a write, and why. Recorded on the write's tick.
145#[derive(Debug, Clone, Default, PartialEq, Eq)]
146#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
147#[cfg_attr(feature = "serde", serde(default))]
148pub struct CommitInfo {
149    pub author: Option<String>,
150    pub message: Option<String>,
151}
152
153/// Options of a level change.
154#[derive(Debug, Clone, Default, PartialEq, Eq)]
155#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
156#[cfg_attr(feature = "serde", serde(default, rename_all = "camelCase"))]
157pub struct LevelChange {
158    /// Create (`Some(true)`) or drop (`Some(false)`) the as-of log indexes.
159    pub as_of_index: Option<bool>,
160    /// Create or drop the index on `quads.t`.
161    pub stamp_index: Option<bool>,
162    /// Allow a downgrade to delete data (the change log, the ticks, the stamp column).
163    pub allow_loss: bool,
164    /// Recorded on the tick of the change.
165    pub author: Option<String>,
166    pub message: Option<String>,
167}
168
169/// Kinds of ticks. Ticks of kind 0 are ordinary writes; the others mark level changes and the
170/// boundaries of the recorded history.
171pub mod kind {
172    /// An atomic write.
173    pub const WRITE: i64 = 0;
174    /// History starts here (upgrade to `log`): the log holds the whole store at this tick.
175    pub const GENESIS: i64 = 1;
176    /// History stopped recording (downgrade from `log`); as-of works up to this tick.
177    pub const FREEZE: i64 = 2;
178    /// History resumed after a freeze: the gap's net changes are recorded on this tick.
179    pub const RESUME: i64 = 3;
180    /// History was deleted (downgrade with `allow_loss`); nothing before it is queryable.
181    pub const DROPPED: i64 = 4;
182    /// Another level change (stamping started or stopped, indexes changed).
183    pub const LEVEL: i64 = 5;
184    /// A purge: quads removed from the store and from the history.
185    pub const PURGE: i64 = 6;
186
187    pub fn name(k: i64) -> &'static str {
188        match k {
189            WRITE => "write",
190            GENESIS => "genesis",
191            FREEZE => "freeze",
192            RESUME => "resume",
193            DROPPED => "dropped",
194            LEVEL => "level",
195            PURGE => "purge",
196            _ => "unknown",
197        }
198    }
199}
200
201/// SQL: the current tick.
202pub const CURRENT_TICK: &str = "(SELECT max(t) FROM ticks)";
203
204/// SQL: now, in seconds since the epoch (works on every SQLite version and on D1).
205pub const NOW: &str = "((julianday('now') - 2440587.5) * 86400.0)";
206
207/// The statement opening a new tick.
208pub fn tick_statement(k: i64, author: Option<&str>, message: Option<&str>) -> Statement {
209    Statement::new(format!(
210        "INSERT INTO ticks(t, time, kind, author, message) SELECT coalesce(max(t), 0) + 1, {NOW}, {k}, {}, {} FROM ticks",
211        sql_opt_str(author),
212        sql_opt_str(message)
213    ))
214}
215
216/// Does this statement write `quads`? Every writer of the asserted quads spells its statement
217/// with one of these two prefixes (see `writer` and `update`); a test holds them to it.
218pub fn writes_quads(s: &Statement) -> bool {
219    s.sql.starts_with("INSERT OR IGNORE INTO quads(") || s.sql.starts_with("DELETE FROM quads ")
220}
221
222fn opens_own_tick(r: &Request) -> bool {
223    r.statements
224        .iter()
225        .take(PREFIX_LEN)
226        .any(|s| s.sql.starts_with("INSERT INTO ticks("))
227}
228
229/// The statements [`prepare`] puts in front of a write: the terms of the tick's time, author and
230/// message, then the tick itself.
231pub const PREFIX_LEN: usize = 2;
232
233/// The statements opening a write tick: the tick's time (read from the host's clock), author and
234/// message become terms, so the history can be queried as RDF (`<oxilite:history>`), and the
235/// tick row records their ids.
236pub fn write_tick_statements(info: &CommitInfo) -> Vec<Statement> {
237    let now = oxsdatatypes::DateTime::now().to_string();
238    let xsd_dt = "http://www.w3.org/2001/XMLSchema#dateTime";
239    let secs = crate::encoding::timestamp(&now, xsd_dt).unwrap_or_default();
240    let mut rows = crate::encoding::EncodedRows::default();
241    let time = oxrdf::Literal::new_typed_literal(now, oxrdf::NamedNode::new_unchecked(xsd_dt));
242    let time_id = rows.term(time.as_ref().into());
243    let mut lit = |v: &Option<String>| {
244        v.as_deref().map_or_else(
245            || "NULL".to_owned(),
246            |v| {
247                rows.term(oxrdf::LiteralRef::new_simple_literal(v).into())
248                    .to_string()
249            },
250        )
251    };
252    let (author_id, message_id) = (lit(&info.author), lit(&info.message));
253    let mut out = crate::writer::term_statements(&rows, &Capabilities::native());
254    out.truncate(1);
255    out.push(Statement::new(format!(
256        "INSERT INTO ticks(t, time, kind, author, message, time_id, author_id, message_id) \
257         SELECT coalesce(max(t), 0) + 1, {}, {}, {}, {}, {time_id}, {author_id}, {message_id} FROM ticks",
258        crate::sql::sql_f64(secs),
259        kind::WRITE,
260        sql_opt_str(info.author.as_deref()),
261        sql_opt_str(info.message.as_deref()),
262    )));
263    out
264}
265
266/// The request with a tick opened first, when the level requires one and the request writes
267/// quads. `None`: run the request unchanged. The first [`PREFIX_LEN`] results belong to the
268/// tick ([`strip`]).
269pub fn prepare(request: &Request, level: Versioning, info: &CommitInfo) -> Option<Request> {
270    if level < Versioning::Stamped
271        || request.mode != Mode::Atomic
272        || opens_own_tick(request)
273        || !request.statements.iter().any(writes_quads)
274    {
275        return None;
276    }
277    let mut statements = write_tick_statements(info);
278    debug_assert_eq!(statements.len(), PREFIX_LEN);
279    statements.extend(request.statements.iter().cloned());
280    Some(Request {
281        statements,
282        mode: request.mode,
283    })
284}
285
286/// Removes the results of the statements added by [`prepare`].
287pub fn strip(mut response: Response) -> Response {
288    response.drain(..PREFIX_LEN.min(response.len()));
289    response
290}
291
292/// The capabilities a store at `level` hands to its jobs: writers stamp quads, and one
293/// statement of every request is reserved for the tick.
294pub fn effective_caps(caps: &Capabilities, level: Versioning) -> Capabilities {
295    let mut c = caps.clone();
296    c.versioning = level;
297    if level >= Versioning::Stamped {
298        c.max_statements = c.max_statements.saturating_sub(PREFIX_LEN).max(1);
299    }
300    c
301}
302
303/// A job whose atomic writes open a tick (for hosts that run requests themselves, such as the
304/// JavaScript drivers).
305pub struct Versioned<J> {
306    inner: J,
307    level: Versioning,
308    info: CommitInfo,
309    pending: bool,
310}
311
312impl<J> Versioned<J> {
313    pub fn new(inner: J, level: Versioning, info: CommitInfo) -> Self {
314        Self {
315            inner,
316            level,
317            info,
318            pending: false,
319        }
320    }
321}
322
323impl<J: Job> Job for Versioned<J> {
324    type Output = J::Output;
325    fn step(&mut self, response: Option<Response>) -> Result<Step<J::Output>> {
326        let response = if std::mem::take(&mut self.pending) {
327            response.map(strip)
328        } else {
329            response
330        };
331        match self.inner.step(response)? {
332            Step::Execute(r) => match prepare(&r, self.level, &self.info) {
333                Some(p) => {
334                    self.pending = true;
335                    Ok(Step::Execute(p))
336                }
337                None => Ok(Step::Execute(r)),
338            },
339            done => Ok(done),
340        }
341    }
342}
343
344/// A backend adapter that opens a tick for every write (one per interactive transaction), and
345/// hands out capabilities matching the store's level. Stores wrap their backend in it.
346pub struct VersionedBackend<B> {
347    inner: B,
348    caps: [Capabilities; 3],
349    level: AtomicU8,
350    info: RwLock<CommitInfo>,
351    in_transaction: AtomicBool,
352    tick_open: AtomicBool,
353}
354
355impl<B> VersionedBackend<B> {
356    pub fn new(inner: B, caps: &Capabilities, level: Versioning) -> Self {
357        Self {
358            inner,
359            caps: Versioning::ALL.map(|l| effective_caps(caps, l)),
360            level: AtomicU8::new(level.index()),
361            info: RwLock::new(CommitInfo::default()),
362            in_transaction: AtomicBool::new(false),
363            tick_open: AtomicBool::new(false),
364        }
365    }
366
367    pub fn inner(&self) -> &B {
368        &self.inner
369    }
370
371    pub fn level(&self) -> Versioning {
372        Versioning::from_index(self.level.load(Ordering::SeqCst))
373    }
374
375    pub fn set_level(&self, level: Versioning) {
376        self.level.store(level.index(), Ordering::SeqCst);
377    }
378
379    /// Author and message recorded on the ticks of later writes (until changed).
380    pub fn set_commit_info(&self, info: CommitInfo) {
381        *self
382            .info
383            .write()
384            .unwrap_or_else(std::sync::PoisonError::into_inner) = info;
385    }
386
387    pub fn commit_info(&self) -> CommitInfo {
388        self.info
389            .read()
390            .unwrap_or_else(std::sync::PoisonError::into_inner)
391            .clone()
392    }
393
394    fn prepared(&self, request: &Request) -> Option<Request> {
395        if self.in_transaction.load(Ordering::SeqCst) && self.tick_open.load(Ordering::SeqCst) {
396            return None;
397        }
398        let p = prepare(request, self.level(), &self.commit_info())?;
399        if self.in_transaction.load(Ordering::SeqCst) {
400            self.tick_open.store(true, Ordering::SeqCst);
401        }
402        Some(p)
403    }
404}
405
406impl<B: SyncBackend> SyncBackend for VersionedBackend<B> {
407    fn execute(&self, request: &Request) -> Result<Response> {
408        match self.prepared(request) {
409            Some(p) => self.inner.execute(&p).map(strip),
410            None => self.inner.execute(request),
411        }
412    }
413    fn capabilities(&self) -> &Capabilities {
414        &self.caps[usize::from(self.level.load(Ordering::SeqCst))]
415    }
416    fn begin(&self) -> Result<()> {
417        self.inner.begin()?;
418        self.in_transaction.store(true, Ordering::SeqCst);
419        self.tick_open.store(false, Ordering::SeqCst);
420        Ok(())
421    }
422    fn commit(&self) -> Result<()> {
423        self.in_transaction.store(false, Ordering::SeqCst);
424        self.tick_open.store(false, Ordering::SeqCst);
425        self.inner.commit()
426    }
427    fn rollback(&self) -> Result<()> {
428        self.in_transaction.store(false, Ordering::SeqCst);
429        self.tick_open.store(false, Ordering::SeqCst);
430        self.inner.rollback()
431    }
432}
433
434impl<B: AsyncBackend> AsyncBackend for VersionedBackend<B> {
435    async fn execute(&self, request: &Request) -> Result<Response> {
436        match self.prepared(request) {
437            Some(p) => self.inner.execute(&p).await.map(strip),
438            None => self.inner.execute(request).await,
439        }
440    }
441    fn capabilities(&self) -> &Capabilities {
442        &self.caps[usize::from(self.level.load(Ordering::SeqCst))]
443    }
444}
445
446// ------------------------------------------------------------------------------------ schema
447
448fn meta(entries: &[(&str, &str)]) -> Statement {
449    Statement::new(format!(
450        "INSERT OR REPLACE INTO oxilite_meta(key, value) VALUES {}",
451        entries
452            .iter()
453            .map(|(k, v)| format!("({}, {})", sql_str(k), sql_str(v)))
454            .collect::<Vec<_>>()
455            .join(", ")
456    ))
457}
458
459const PURGING: &str = "EXISTS (SELECT 1 FROM oxilite_meta WHERE key = 'purging')";
460
461/// The clock: `ticks`, immutable except during a purge.
462fn clock_ddl() -> Vec<Statement> {
463    [
464        // `time_id`, `author_id`, `message_id`: the terms of a write tick (see `write_tick_statements`).
465        "CREATE TABLE IF NOT EXISTS ticks (t INTEGER PRIMARY KEY, time REAL NOT NULL, kind INTEGER NOT NULL DEFAULT 0, author TEXT, message TEXT, time_id INTEGER, author_id INTEGER, message_id INTEGER) STRICT".to_owned(),
466        // Boundaries of the recorded history, found without scanning the write ticks.
467        "CREATE INDEX IF NOT EXISTS ticks_boundary ON ticks(t) WHERE kind > 0".to_owned(),
468        format!("CREATE TRIGGER IF NOT EXISTS ticks_immutable_update BEFORE UPDATE ON ticks WHEN NOT {PURGING} BEGIN SELECT RAISE(ABORT, 'oxilite: history is immutable'); END"),
469        format!("CREATE TRIGGER IF NOT EXISTS ticks_immutable_delete BEFORE DELETE ON ticks WHEN NOT {PURGING} BEGIN SELECT RAISE(ABORT, 'oxilite: history is immutable'); END"),
470    ]
471    .into_iter()
472    .map(Statement::new)
473    .collect()
474}
475
476/// The change log and its triggers. A change of the tick being written may be undone within
477/// the same tick (a quad added then removed leaves no trace); earlier ticks are immutable.
478fn log_tables_ddl() -> Vec<Statement> {
479    let past = format!("OLD.tx < {CURRENT_TICK} AND NOT {PURGING}");
480    [
481        "CREATE TABLE IF NOT EXISTS quad_log (s INTEGER NOT NULL, p INTEGER NOT NULL, o INTEGER NOT NULL, g INTEGER NOT NULL, tx INTEGER NOT NULL, op INTEGER NOT NULL, PRIMARY KEY (s, p, o, g, tx)) WITHOUT ROWID, STRICT".to_owned(),
482        // The changes of one commit (log listings, diffs, change feeds).
483        "CREATE INDEX IF NOT EXISTS quad_log_tx ON quad_log(tx)".to_owned(),
484        "CREATE TABLE IF NOT EXISTS commits (tx INTEGER PRIMARY KEY) STRICT".to_owned(),
485        format!("CREATE TRIGGER IF NOT EXISTS quad_log_immutable_update BEFORE UPDATE ON quad_log WHEN {past} BEGIN SELECT RAISE(ABORT, 'oxilite: history is immutable'); END"),
486        format!("CREATE TRIGGER IF NOT EXISTS quad_log_immutable_delete BEFORE DELETE ON quad_log WHEN {past} BEGIN SELECT RAISE(ABORT, 'oxilite: history is immutable'); END"),
487        format!("CREATE TRIGGER IF NOT EXISTS commits_immutable_delete BEFORE DELETE ON commits WHEN NOT {PURGING} BEGIN SELECT RAISE(ABORT, 'oxilite: history is immutable'); END"),
488    ]
489    .into_iter()
490    .map(Statement::new)
491    .collect()
492}
493
494/// The triggers recording every effective change of `quads` in the log.
495fn capture_ddl() -> Vec<Statement> {
496    let same = |r: &str, x: &str| {
497        format!("{r}.s = {x}.s AND {r}.p = {x}.p AND {r}.o = {x}.o AND {r}.g = {x}.g")
498    };
499    let ins_same = same("l", "NEW");
500    let del_same = same("l", "OLD");
501    // SQLite forbids aliases on the target of a DELETE inside a trigger (D1 enforces it).
502    let plain = |x: &str| format!("s = {x}.s AND p = {x}.p AND o = {x}.o AND g = {x}.g");
503    let (ins_plain, del_plain) = (plain("NEW"), plain("OLD"));
504    vec![
505        Statement::new(format!(
506            "CREATE TRIGGER IF NOT EXISTS quads_log_insert AFTER INSERT ON quads WHEN NOT {PURGING} BEGIN \
507             INSERT INTO quad_log(s, p, o, g, tx, op) SELECT NEW.s, NEW.p, NEW.o, NEW.g, {CURRENT_TICK}, 1 \
508               WHERE NOT EXISTS (SELECT 1 FROM quad_log l WHERE {ins_same} AND l.tx = {CURRENT_TICK} AND l.op = 0); \
509             DELETE FROM quad_log WHERE {ins_plain} AND tx = {CURRENT_TICK} AND op = 0; \
510             INSERT OR IGNORE INTO commits(tx) SELECT {CURRENT_TICK}; \
511             END"
512        )),
513        Statement::new(format!(
514            "CREATE TRIGGER IF NOT EXISTS quads_log_delete AFTER DELETE ON quads WHEN NOT {PURGING} BEGIN \
515             INSERT INTO quad_log(s, p, o, g, tx, op) SELECT OLD.s, OLD.p, OLD.o, OLD.g, {CURRENT_TICK}, 0 \
516               WHERE NOT EXISTS (SELECT 1 FROM quad_log l WHERE {del_same} AND l.tx = {CURRENT_TICK} AND l.op = 1); \
517             DELETE FROM quad_log WHERE {del_plain} AND tx = {CURRENT_TICK} AND op = 1; \
518             INSERT OR IGNORE INTO commits(tx) SELECT {CURRENT_TICK}; \
519             END"
520        )),
521    ]
522}
523
524fn as_of_index_ddl(on: bool) -> Vec<Statement> {
525    if on {
526        vec![
527            "CREATE INDEX IF NOT EXISTS quad_log_posg ON quad_log(p, o, s, g, tx)".into(),
528            "CREATE INDEX IF NOT EXISTS quad_log_ospg ON quad_log(o, s, p, g, tx)".into(),
529        ]
530    } else {
531        vec![
532            "DROP INDEX IF EXISTS quad_log_posg".into(),
533            "DROP INDEX IF EXISTS quad_log_ospg".into(),
534        ]
535    }
536}
537
538fn stamp_index_ddl(on: bool) -> Statement {
539    if on {
540        "CREATE INDEX IF NOT EXISTS quads_t ON quads(t)".into()
541    } else {
542        "DROP INDEX IF EXISTS quads_t".into()
543    }
544}
545
546/// The statements moving a store from `from` to level `to` (one level at a time, in order), then
547/// applying the index options of `change`. One atomic request; also the body of a D1 migration.
548///
549/// Upgrades never lose data; the upgrade to `log` records the whole store as its genesis
550/// commit (or, after a freeze, the net changes of the gap). Downgrades keep data unless
551/// `allow_loss`: the log freezes, and the clock stops.
552pub fn change_statements(
553    from: &VersionState,
554    to: Versioning,
555    change: &LevelChange,
556) -> Result<Vec<Statement>> {
557    let mut state = *from;
558    let mut out = Vec::new();
559    let author = change.author.as_deref();
560    let msg = |default: &str| Some(change.message.clone().unwrap_or_else(|| default.to_owned()));
561    while state.level < to {
562        match state.level {
563            Versioning::Off => {
564                out.extend(clock_ddl());
565                if !state.stamp_column {
566                    out.push("ALTER TABLE quads ADD COLUMN t INTEGER NOT NULL DEFAULT 0".into());
567                }
568                out.push(tick_statement(
569                    kind::LEVEL,
570                    author,
571                    msg("versioning: stamped").as_deref(),
572                ));
573                out.push(meta(&[("versioning", "stamped"), ("stamp_column", "1")]));
574                state.level = Versioning::Stamped;
575                state.stamp_column = true;
576            }
577            Versioning::Stamped => {
578                out.extend(log_tables_ddl());
579                if state.history == History::Frozen {
580                    // The gap's net changes, against the state the log froze at.
581                    let freeze =
582                        format!("(SELECT max(t) FROM ticks WHERE kind = {})", kind::FREEZE);
583                    let frozen = as_of_sql(&freeze);
584                    out.push(tick_statement(
585                        kind::RESUME,
586                        author,
587                        msg("history resumed").as_deref(),
588                    ));
589                    out.push(Statement::new(format!(
590                        "INSERT INTO quad_log(s, p, o, g, tx, op) SELECT f.s, f.p, f.o, f.g, {CURRENT_TICK}, 0 FROM {frozen} f \
591                         WHERE NOT EXISTS (SELECT 1 FROM quads q WHERE q.s = f.s AND q.p = f.p AND q.o = f.o AND q.g = f.g)"
592                    )));
593                    out.push(Statement::new(format!(
594                        "INSERT INTO quad_log(s, p, o, g, tx, op) SELECT q.s, q.p, q.o, q.g, {CURRENT_TICK}, 1 FROM quads q \
595                         WHERE NOT EXISTS (SELECT 1 FROM {frozen} f WHERE q.s = f.s AND q.p = f.p AND q.o = f.o AND q.g = f.g)"
596                    )));
597                } else {
598                    out.push(tick_statement(
599                        kind::GENESIS,
600                        author,
601                        msg("history starts").as_deref(),
602                    ));
603                    out.push(Statement::new(format!(
604                        "INSERT INTO quad_log(s, p, o, g, tx, op) SELECT s, p, o, g, {CURRENT_TICK}, 1 FROM quads"
605                    )));
606                }
607                out.push(format!("INSERT OR IGNORE INTO commits(tx) SELECT {CURRENT_TICK}").into());
608                out.extend(capture_ddl());
609                out.push(meta(&[("versioning", "log"), ("history", "live")]));
610                state.level = Versioning::Log;
611                state.history = History::Live;
612            }
613            Versioning::Log => unreachable!(),
614        }
615    }
616    while state.level > to {
617        match state.level {
618            Versioning::Log => {
619                out.push("DROP TRIGGER IF EXISTS quads_log_insert".into());
620                out.push("DROP TRIGGER IF EXISTS quads_log_delete".into());
621                if change.allow_loss {
622                    out.push("DROP TABLE IF EXISTS quad_log".into());
623                    out.push("DROP TABLE IF EXISTS commits".into());
624                    out.push(tick_statement(
625                        kind::DROPPED,
626                        author,
627                        msg("history deleted").as_deref(),
628                    ));
629                    out.push(meta(&[
630                        ("versioning", "stamped"),
631                        ("history", "none"),
632                        ("as_of_index", "0"),
633                    ]));
634                    state.history = History::None;
635                    state.as_of_index = false;
636                } else {
637                    out.push(tick_statement(
638                        kind::FREEZE,
639                        author,
640                        msg("history frozen").as_deref(),
641                    ));
642                    out.push(meta(&[("versioning", "stamped"), ("history", "frozen")]));
643                    state.history = History::Frozen;
644                }
645                state.level = Versioning::Stamped;
646            }
647            Versioning::Stamped => {
648                if change.allow_loss {
649                    if state.history != History::None {
650                        return Err(Error::Other(
651                            "the frozen history needs the clock: delete it first (lower the level to `stamped` with allow_loss from `log`)".into(),
652                        ));
653                    }
654                    out.push("DROP INDEX IF EXISTS quads_t".into());
655                    out.push("ALTER TABLE quads DROP COLUMN t".into());
656                    out.push("DROP TABLE IF EXISTS ticks".into());
657                    out.push(meta(&[
658                        ("versioning", "off"),
659                        ("stamp_column", "0"),
660                        ("stamp_index", "0"),
661                    ]));
662                    state.stamp_column = false;
663                    state.stamp_index = false;
664                } else {
665                    out.push(tick_statement(
666                        kind::LEVEL,
667                        author,
668                        msg("versioning: off").as_deref(),
669                    ));
670                    out.push(meta(&[("versioning", "off")]));
671                }
672                state.level = Versioning::Off;
673            }
674            Versioning::Off => unreachable!(),
675        }
676    }
677    if let Some(on) = change.as_of_index {
678        if on && state.history == History::None {
679            return Err(Error::Other(
680                "the as-of index needs a change log (level `log`)".into(),
681            ));
682        }
683        if on != state.as_of_index {
684            out.extend(as_of_index_ddl(on));
685            out.push(meta(&[("as_of_index", if on { "1" } else { "0" })]));
686        }
687    }
688    if let Some(on) = change.stamp_index {
689        if on && !state.stamp_column {
690            return Err(Error::Other(
691                "the stamp index needs the store clock (level `stamped`)".into(),
692            ));
693        }
694        if on != state.stamp_index {
695            out.push(stamp_index_ddl(on));
696            out.push(meta(&[("stamp_index", if on { "1" } else { "0" })]));
697        }
698    }
699    Ok(out)
700}
701
702// ------------------------------------------------------------------------------------- as-of
703
704/// SQL: the `(s, p, o, g)` table of the store as it was at tick `t` (a SQL expression): the
705/// quads whose latest logged change at or before `t` is an addition.
706pub fn as_of_sql(t: &str) -> String {
707    format!(
708        "(SELECT l.s AS s, l.p AS p, l.o AS o, l.g AS g FROM quad_log l WHERE l.tx <= {t} AND l.op = 1 \
709         AND NOT EXISTS (SELECT 1 FROM quad_log r WHERE r.s = l.s AND r.p = l.p AND r.o = l.o AND r.g = l.g AND r.tx > l.tx AND r.tx <= {t}))"
710    )
711}
712
713/// A version of the store, as users name it.
714#[derive(Debug, Clone, PartialEq)]
715pub enum VersionRef {
716    /// The latest tick (`HEAD`, `main`).
717    Head,
718    /// `n` commits before the latest (`HEAD~n`, `main~n`, `~n`).
719    Back(u64),
720    /// A tick number (`#42`, `42`).
721    Tick(i64),
722    /// The state at an instant (`@2026-09-01T12:00:00Z`, `HEAD@…`, `main@…`), in seconds since
723    /// the epoch.
724    Time(f64),
725}
726
727impl FromStr for VersionRef {
728    type Err = Error;
729    fn from_str(s: &str) -> Result<Self> {
730        let s = s.trim();
731        let bad = || {
732            Error::Other(format!(
733                "invalid version `{s}` (expected HEAD, HEAD~n, a tick number such as #42, or @<xsd:dateTime>)"
734            ))
735        };
736        let rest = s
737            .strip_prefix("HEAD")
738            .or_else(|| s.strip_prefix("main"))
739            .unwrap_or(s);
740        if rest.is_empty() {
741            return Ok(Self::Head);
742        }
743        if let Some(n) = rest.strip_prefix('~') {
744            return n.parse().map(Self::Back).map_err(|_| bad());
745        }
746        if let Some(t) = rest.strip_prefix('@') {
747            let dt = if t.len() == 10 {
748                "http://www.w3.org/2001/XMLSchema#date"
749            } else {
750                "http://www.w3.org/2001/XMLSchema#dateTime"
751            };
752            return crate::encoding::timestamp(t, dt)
753                .map(Self::Time)
754                .ok_or_else(bad);
755        }
756        let n = rest.strip_prefix('#').unwrap_or(rest);
757        n.parse().map(Self::Tick).map_err(|_| bad())
758    }
759}
760
761impl fmt::Display for VersionRef {
762    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
763        match self {
764            Self::Head => f.write_str("HEAD"),
765            Self::Back(n) => write!(f, "HEAD~{n}"),
766            Self::Tick(t) => write!(f, "#{t}"),
767            Self::Time(t) => write!(f, "@{}", format_time(*t)),
768        }
769    }
770}
771
772fn ref_sql(r: &VersionRef) -> String {
773    match r {
774        VersionRef::Head => "SELECT max(t) FROM ticks".to_owned(),
775        VersionRef::Back(n) => {
776            format!("SELECT tx FROM commits ORDER BY tx DESC LIMIT 1 OFFSET {n}")
777        }
778        VersionRef::Tick(t) => format!("SELECT t FROM ticks WHERE t = {t}"),
779        VersionRef::Time(x) => format!(
780            "SELECT max(t) FROM ticks WHERE time <= {}",
781            crate::sql::sql_f64(*x)
782        ),
783    }
784}
785
786/// Resolves version references to ticks, checking that each lies inside the recorded history.
787pub fn resolve_job(refs: Vec<VersionRef>, state: VersionState) -> ResolveJob {
788    ResolveJob {
789        refs,
790        state,
791        sent: false,
792    }
793}
794
795pub struct ResolveJob {
796    refs: Vec<VersionRef>,
797    state: VersionState,
798    sent: bool,
799}
800
801impl Job for ResolveJob {
802    type Output = Vec<i64>;
803    fn step(&mut self, response: Option<Response>) -> Result<Step<Vec<i64>>> {
804        if self.state.history == History::None {
805            return Err(Error::Other(
806                "this store keeps no history: raise its versioning level to `log` for as-of queries".into(),
807            ));
808        }
809        if self.refs.is_empty() {
810            return Ok(Step::Done(Vec::new()));
811        }
812        if !std::mem::replace(&mut self.sent, true) {
813            let b = format!(
814                "kind IN ({}, {}, {}, {})",
815                kind::GENESIS,
816                kind::FREEZE,
817                kind::RESUME,
818                kind::DROPPED
819            );
820            return Ok(Step::Execute(Request::read(
821                self.refs
822                    .iter()
823                    .map(|r| {
824                        Statement::new(format!(
825                            "WITH r(t) AS ({}) SELECT r.t, (SELECT b.t FROM ticks b WHERE {b} AND b.t <= r.t ORDER BY b.t DESC LIMIT 1), \
826                             (SELECT b.kind FROM ticks b WHERE {b} AND b.t <= r.t ORDER BY b.t DESC LIMIT 1) FROM r",
827                            ref_sql(r)
828                        ))
829                    })
830                    .collect(),
831            )));
832        }
833        let response = response.unwrap_or_default();
834        expect_len(&response, self.refs.len())?;
835        let mut out = Vec::with_capacity(self.refs.len());
836        for (r, rs) in self.refs.iter().zip(&response) {
837            let row = rs.rows.first();
838            let t = row.and_then(|row| col(row, 0).ok()?.as_i64());
839            let Some(t) = t else {
840                return Err(Error::Other(format!("unknown version {r}")));
841            };
842            let row = row.expect("row with a tick");
843            let boundary = col(row, 1)?.as_i64();
844            let k = col(row, 2)?.as_i64();
845            match (boundary, k) {
846                (Some(_), Some(kind::GENESIS | kind::RESUME)) => {}
847                (Some(b), Some(kind::FREEZE)) if t == b => {}
848                (Some(b), Some(kind::FREEZE)) => {
849                    return Err(Error::Other(format!(
850                        "version {r} (#{t}) falls in a gap: history was frozen at #{b} and not recorded after it"
851                    )))
852                }
853                _ => {
854                    return Err(Error::Other(format!(
855                        "version {r} (#{t}) is before the recorded history"
856                    )))
857                }
858            }
859            out.push(t);
860        }
861        Ok(Step::Done(out))
862    }
863}
864
865/// `YYYY-MM-DDThh:mm:ss.sssZ` for seconds since the epoch.
866pub fn format_time(secs: f64) -> String {
867    let ms = (secs * 1000.0).round() as i64;
868    let (days, rem) = (ms.div_euclid(86_400_000), ms.rem_euclid(86_400_000));
869    // Civil from days (Howard Hinnant's algorithm).
870    let z = days + 719_468;
871    let era = z.div_euclid(146_097);
872    let doe = z.rem_euclid(146_097);
873    let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
874    let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
875    let mp = (5 * doy + 2) / 153;
876    let d = doy - (153 * mp + 2) / 5 + 1;
877    let m = if mp < 10 { mp + 3 } else { mp - 9 };
878    let y = yoe + era * 400 + i64::from(m <= 2);
879    format!(
880        "{y:04}-{m:02}-{d:02}T{:02}:{:02}:{:02}.{:03}Z",
881        rem / 3_600_000,
882        rem / 60_000 % 60,
883        rem / 1000 % 60,
884        rem % 1000
885    )
886}
887
888// ------------------------------------------------------------------------------- history jobs
889
890fn id_col(caps: &Capabilities, c: &str) -> String {
891    if caps.int64_as_text {
892        format!("CAST({c} AS TEXT)")
893    } else {
894        c.to_owned()
895    }
896}
897
898fn opt_i64(v: &crate::sql::SqlValue) -> Option<i64> {
899    v.as_i64()
900}
901
902fn opt_string(v: &crate::sql::SqlValue) -> Option<String> {
903    v.clone().into_string()
904}
905
906/// The versioning state of a store and where its clock and history stand.
907#[derive(Debug, Clone, PartialEq, Default)]
908#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
909#[cfg_attr(feature = "serde", serde(rename_all = "camelCase"))]
910pub struct VersionStatus {
911    #[cfg_attr(feature = "serde", serde(flatten))]
912    pub state: VersionState,
913    /// The latest tick.
914    pub head: Option<i64>,
915    /// Wall time of the latest tick (seconds since the epoch).
916    pub head_time: Option<f64>,
917    /// Where the recorded history (re)starts: the latest genesis or resume tick.
918    pub genesis: Option<i64>,
919    /// The latest freeze, while the history is frozen.
920    pub frozen_at: Option<i64>,
921    /// Commits in the change log.
922    pub commits: Option<u64>,
923}
924
925/// Reads the [`VersionStatus`].
926pub fn status_job(state: VersionState) -> crate::job::OneShot<VersionStatus> {
927    let mut stmts = Vec::new();
928    if state.has_ticks() {
929        stmts.push(Statement::new(format!(
930            "SELECT max(t), (SELECT time FROM ticks WHERE t = (SELECT max(t) FROM ticks)), \
931             (SELECT max(t) FROM ticks WHERE kind IN ({}, {})), (SELECT max(t) FROM ticks WHERE kind = {}) FROM ticks",
932            kind::GENESIS,
933            kind::RESUME,
934            kind::FREEZE
935        )));
936    }
937    if state.history != History::None {
938        stmts.push("SELECT count(*) FROM commits".into());
939    }
940    crate::job::OneShot::new(Request::read(stmts), move |r: Response| {
941        let mut s = VersionStatus {
942            state,
943            ..Default::default()
944        };
945        let mut rs = r.iter();
946        if state.has_ticks() {
947            if let Some(row) = rs.next().and_then(|x| x.rows.first()) {
948                s.head = opt_i64(col(row, 0)?);
949                s.head_time = col(row, 1)?.as_f64();
950                s.genesis = opt_i64(col(row, 2)?);
951                let freeze = opt_i64(col(row, 3)?);
952                if state.history == History::Frozen {
953                    s.frozen_at = freeze;
954                }
955            }
956        }
957        if state.history != History::None {
958            if let Some(row) = rs.next().and_then(|x| x.rows.first()) {
959                s.commits = opt_i64(col(row, 0)?).map(|n| n as u64);
960            }
961        }
962        Ok(s)
963    })
964}
965
966/// One entry of the history: a commit, or a level change.
967#[derive(Debug, Clone, PartialEq)]
968#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
969#[cfg_attr(feature = "serde", serde(rename_all = "camelCase"))]
970pub struct CommitRecord {
971    pub tick: i64,
972    /// Seconds since the epoch.
973    pub time: f64,
974    /// `write`, `genesis`, `freeze`, `resume`, `dropped`, `level` or `purge`.
975    pub kind: String,
976    pub author: Option<String>,
977    pub message: Option<String>,
978    /// Quads added and removed (known while a change log exists).
979    pub added: Option<u64>,
980    pub removed: Option<u64>,
981}
982
983/// The latest `limit` entries of the history, newest first.
984pub fn log_job(
985    state: VersionState,
986    limit: usize,
987) -> Result<crate::job::OneShot<Vec<CommitRecord>>> {
988    if !state.has_ticks() {
989        return Err(Error::Other(
990            "this store keeps no history: raise its versioning level to `stamped` or `log`".into(),
991        ));
992    }
993    let sql = if state.history == History::None {
994        format!("SELECT t, time, kind, author, message, NULL, NULL FROM ticks ORDER BY t DESC LIMIT {limit}")
995    } else {
996        let filter = if state.level == Versioning::Log {
997            "WHERE k.kind <> 0 OR EXISTS (SELECT 1 FROM commits c WHERE c.tx = k.t)"
998        } else {
999            ""
1000        };
1001        format!(
1002            "SELECT k.t, k.time, k.kind, k.author, k.message, \
1003             (SELECT count(*) FROM quad_log l WHERE l.tx = k.t AND l.op = 1), \
1004             (SELECT count(*) FROM quad_log l WHERE l.tx = k.t AND l.op = 0) \
1005             FROM ticks k {filter} ORDER BY k.t DESC LIMIT {limit}"
1006        )
1007    };
1008    Ok(crate::job::OneShot::new(
1009        Request::read(vec![Statement::new(sql)]),
1010        |r: Response| {
1011            expect_len(&r, 1)?;
1012            r[0].rows
1013                .iter()
1014                .map(|row| {
1015                    Ok(CommitRecord {
1016                        tick: opt_i64(col(row, 0)?).unwrap_or_default(),
1017                        time: col(row, 1)?.as_f64().unwrap_or_default(),
1018                        kind: kind::name(opt_i64(col(row, 2)?).unwrap_or_default()).to_owned(),
1019                        author: opt_string(col(row, 3)?),
1020                        message: opt_string(col(row, 4)?),
1021                        added: opt_i64(col(row, 5)?).map(|n| n as u64),
1022                        removed: opt_i64(col(row, 6)?).map(|n| n as u64),
1023                    })
1024                })
1025                .collect()
1026        },
1027    ))
1028}
1029
1030/// A quad added or removed at a tick.
1031#[derive(Debug, Clone, PartialEq)]
1032pub struct Change {
1033    pub tick: i64,
1034    pub added: bool,
1035    pub quad: oxrdf::Quad,
1036}
1037
1038/// Runs `sql` (rows `tick, op, s, p, o, g`) and resolves the terms.
1039pub struct ChangeJob {
1040    sql: Option<String>,
1041    caps: Capabilities,
1042    resolver: crate::resolve::TermResolver,
1043    rows: Vec<(i64, bool, [i64; 4])>,
1044    started: bool,
1045}
1046
1047impl ChangeJob {
1048    fn new(sql: String, caps: &Capabilities) -> Self {
1049        Self {
1050            sql: Some(sql),
1051            caps: caps.clone(),
1052            resolver: crate::resolve::TermResolver::default(),
1053            rows: Vec::new(),
1054            started: false,
1055        }
1056    }
1057}
1058
1059impl Job for ChangeJob {
1060    type Output = Vec<Change>;
1061    fn step(&mut self, response: Option<Response>) -> Result<Step<Vec<Change>>> {
1062        if let Some(sql) = self.sql.take() {
1063            return Ok(Step::Execute(Request::read(vec![Statement::new(sql)])));
1064        }
1065        let response = response.unwrap_or_default();
1066        if self.started {
1067            self.resolver.absorb(response)?;
1068        } else {
1069            self.started = true;
1070            for rs in response {
1071                for row in rs.rows {
1072                    let ids: Vec<i64> = row
1073                        .iter()
1074                        .filter_map(crate::sql::SqlValue::as_i64)
1075                        .collect();
1076                    let [t, op, s, p, o, g] = ids[..] else {
1077                        return Err(Error::corrupted("bad change row"));
1078                    };
1079                    for id in [s, p, o] {
1080                        self.resolver.want(id);
1081                    }
1082                    if g != crate::encoding::DEFAULT_GRAPH_ID {
1083                        self.resolver.want(g);
1084                    }
1085                    self.rows.push((t, op == 1, [s, p, o, g]));
1086                }
1087            }
1088        }
1089        if let Some(r) = self.resolver.request(&self.caps) {
1090            return Ok(Step::Execute(r));
1091        }
1092        let out = self
1093            .rows
1094            .iter()
1095            .map(|(t, added, [s, p, o, g])| {
1096                let gname = if *g == crate::encoding::DEFAULT_GRAPH_ID {
1097                    oxrdf::GraphName::DefaultGraph
1098                } else {
1099                    crate::encoding::to_graph_name(*g, Some(self.resolver.get(*g)?))?
1100                };
1101                Ok(Change {
1102                    tick: *t,
1103                    added: *added,
1104                    quad: crate::encoding::make_quad(
1105                        self.resolver.get(*s)?,
1106                        self.resolver.get(*p)?,
1107                        self.resolver.get(*o)?,
1108                        gname,
1109                    )?,
1110                })
1111            })
1112            .collect::<Result<_>>()?;
1113        Ok(Step::Done(out))
1114    }
1115}
1116
1117/// The changes recorded after tick `after` up to tick `until` (inclusive), in order. With a
1118/// change log: every addition and removal. With only the clock: the quads present now that were
1119/// added in that range (removals leave no trace).
1120pub fn changes_job(
1121    state: VersionState,
1122    after: i64,
1123    until: Option<i64>,
1124    caps: &Capabilities,
1125) -> Result<ChangeJob> {
1126    let upto = until
1127        .map(|u| format!(" AND {{t}} <= {u}"))
1128        .unwrap_or_default();
1129    let cols = |t: &str, op: &str| {
1130        format!(
1131            "{}, {op}, {}, {}, {}, {}",
1132            id_col(caps, t),
1133            id_col(caps, "s"),
1134            id_col(caps, "p"),
1135            id_col(caps, "o"),
1136            id_col(caps, "g")
1137        )
1138    };
1139    let sql = if state.history != History::None {
1140        format!(
1141            "SELECT {} FROM quad_log WHERE tx > {after}{} ORDER BY tx, op",
1142            cols("tx", "op"),
1143            upto.replace("{t}", "tx")
1144        )
1145    } else if state.stamp_column {
1146        format!(
1147            "SELECT {} FROM quads WHERE t > {after}{} ORDER BY t",
1148            cols("t", "1"),
1149            upto.replace("{t}", "t")
1150        )
1151    } else {
1152        return Err(Error::Other(
1153            "this store keeps no history: raise its versioning level to `stamped` or `log`".into(),
1154        ));
1155    };
1156    Ok(ChangeJob::new(sql, caps))
1157}
1158
1159/// The net difference between two versions (ticks): quads added and removed from `from` to
1160/// `to`, each reported at tick `to`.
1161pub fn diff_job(state: VersionState, from: i64, to: i64, caps: &Capabilities) -> Result<ChangeJob> {
1162    if state.history == History::None {
1163        return Err(Error::Other(
1164            "this store keeps no history: raise its versioning level to `log` to compare versions"
1165                .into(),
1166        ));
1167    }
1168    let (a, b) = (as_of_sql(&from.to_string()), as_of_sql(&to.to_string()));
1169    let cols = |x: &str| {
1170        format!(
1171            "{}, {}, {}, {}",
1172            id_col(caps, &format!("{x}.s")),
1173            id_col(caps, &format!("{x}.p")),
1174            id_col(caps, &format!("{x}.o")),
1175            id_col(caps, &format!("{x}.g"))
1176        )
1177    };
1178    let not_in = |x: &str, other: &str| {
1179        format!("NOT EXISTS (SELECT 1 FROM {other} y WHERE y.s = {x}.s AND y.p = {x}.p AND y.o = {x}.o AND y.g = {x}.g)")
1180    };
1181    let sql = format!(
1182        "SELECT {to}, 1, {} FROM {b} x WHERE {} UNION ALL SELECT {to}, 0, {} FROM {a} x WHERE {}",
1183        cols("x"),
1184        not_in("x", &a),
1185        cols("x"),
1186        not_in("x", &b)
1187    );
1188    Ok(ChangeJob::new(sql, caps))
1189}
1190
1191/// Removes the quads matching the pattern (term ids; `None` matches anything) from the store and
1192/// from the whole history, recording a purge tick that names no removed content. For erasure
1193/// requests: the only operation that rewrites history.
1194pub fn purge_request(
1195    state: VersionState,
1196    pattern: [Option<i64>; 4],
1197    author: Option<&str>,
1198    reason: Option<&str>,
1199) -> Request {
1200    let w = ["s", "p", "o", "g"]
1201        .iter()
1202        .zip(pattern)
1203        .filter_map(|(c, v)| v.map(|v| format!("{c} = {v}")))
1204        .collect::<Vec<_>>();
1205    let w = if w.is_empty() {
1206        "1".to_owned()
1207    } else {
1208        w.join(" AND ")
1209    };
1210    let mut s = Vec::new();
1211    if state.has_ticks() {
1212        s.push(tick_statement(
1213            kind::PURGE,
1214            author,
1215            Some(reason.unwrap_or("purge")),
1216        ));
1217    }
1218    s.push("INSERT OR REPLACE INTO oxilite_meta(key, value) VALUES ('purging', '1')".into());
1219    s.push(Statement::new(format!("DELETE FROM quads WHERE {w}")));
1220    if state.history != History::None {
1221        s.push(Statement::new(format!("DELETE FROM quad_log WHERE {w}")));
1222    }
1223    s.push("DELETE FROM oxilite_meta WHERE key = 'purging'".into());
1224    Request::atomic(s)
1225}
1226
1227/// Rejects version options a query cannot honour: an unresolved version, a store without a
1228/// change log, and inferences or reasoning, which exist only for the current state.
1229pub fn check_query_options(stats: &crate::Stats, options: &crate::QueryOptions) -> Result<()> {
1230    if options.as_of.is_some() && options.as_of_tick.is_none() {
1231        return Err(Error::Other(
1232            "the as-of version was not resolved before compiling".into(),
1233        ));
1234    }
1235    if options.as_of_tick.is_none() && options.versions.is_empty() {
1236        return Ok(());
1237    }
1238    if stats.version.history == History::None {
1239        return Err(Error::Other(
1240            "this store keeps no history: raise its versioning level to `log` for as-of queries"
1241                .into(),
1242        ));
1243    }
1244    if options.include_inferred || options.reasoning != crate::reason::Reasoning::None {
1245        return Err(Error::Other(
1246            "inferences and reasoning describe the current state only; they cannot be combined with an as-of version".into(),
1247        ));
1248    }
1249    Ok(())
1250}
1251
1252/// The version references a query names: its `as_of` option, then each
1253/// `SERVICE <oxilite:version/REF>` of its pattern (deduplicated, in order).
1254pub fn query_version_refs(query: &spargebra::Query, options: &crate::QueryOptions) -> Vec<String> {
1255    // A version already resolved (by a Cypher statement, say) is not resolved again.
1256    let mut out: Vec<String> = if options.as_of_tick.is_none() {
1257        options.as_of.iter().cloned().collect()
1258    } else {
1259        Vec::new()
1260    };
1261    let pattern = match query {
1262        spargebra::Query::Select { pattern, .. }
1263        | spargebra::Query::Construct { pattern, .. }
1264        | spargebra::Query::Describe { pattern, .. }
1265        | spargebra::Query::Ask { pattern, .. } => pattern,
1266    };
1267    collect_services(pattern, &mut out);
1268    let mut seen = std::collections::HashSet::new();
1269    out.retain(|r| {
1270        r != HISTORY_MARKER && !options.versions.contains_key(r) && seen.insert(r.clone())
1271    });
1272    out
1273}
1274
1275/// Does the query read the history graph? Such queries must compile: the fallback evaluator
1276/// would see an empty named graph.
1277pub fn query_reads_history(query: &spargebra::Query) -> bool {
1278    let mut out = Vec::new();
1279    let pattern = match query {
1280        spargebra::Query::Select { pattern, .. }
1281        | spargebra::Query::Construct { pattern, .. }
1282        | spargebra::Query::Describe { pattern, .. }
1283        | spargebra::Query::Ask { pattern, .. } => pattern,
1284    };
1285    collect_services(pattern, &mut out);
1286    out.iter().any(|r| r == HISTORY_MARKER)
1287}
1288
1289const HISTORY_MARKER: &str = "\u{1}history";
1290
1291fn collect_services(p: &spargebra::algebra::GraphPattern, out: &mut Vec<String>) {
1292    use spargebra::algebra::GraphPattern as G;
1293    if let G::Graph {
1294        name: spargebra::term::NamedNodePattern::NamedNode(n),
1295        ..
1296    } = p
1297    {
1298        if n.as_str() == crate::compiler::HISTORY_GRAPH {
1299            out.push(HISTORY_MARKER.to_owned());
1300        }
1301    }
1302    match p {
1303        G::Service { name, inner, .. } => {
1304            if let spargebra::term::NamedNodePattern::NamedNode(n) = name {
1305                if let Some(r) = n.as_str().strip_prefix(crate::compiler::VERSION_SERVICE) {
1306                    out.push(r.to_owned());
1307                }
1308            }
1309            collect_services(inner, out);
1310        }
1311        G::Join { left, right }
1312        | G::LeftJoin { left, right, .. }
1313        | G::Minus { left, right }
1314        | G::Union { left, right } => {
1315            collect_services(left, out);
1316            collect_services(right, out);
1317        }
1318        G::Filter { inner, .. }
1319        | G::Graph { inner, .. }
1320        | G::Extend { inner, .. }
1321        | G::OrderBy { inner, .. }
1322        | G::Project { inner, .. }
1323        | G::Distinct { inner }
1324        | G::Reduced { inner }
1325        | G::Slice { inner, .. }
1326        | G::Group { inner, .. } => collect_services(inner, out),
1327        _ => {}
1328    }
1329}
1330
1331/// Resolves the versions a query names and records their ticks in the options.
1332pub fn resolve_query_job(
1333    refs: Vec<String>,
1334    state: VersionState,
1335    options: crate::QueryOptions,
1336) -> Result<ResolveQueryJob> {
1337    let parsed = refs
1338        .iter()
1339        .map(|r| r.parse::<VersionRef>())
1340        .collect::<Result<Vec<_>>>()?;
1341    Ok(ResolveQueryJob {
1342        inner: resolve_job(parsed, state),
1343        refs,
1344        options: Some(options),
1345    })
1346}
1347
1348pub struct ResolveQueryJob {
1349    inner: ResolveJob,
1350    refs: Vec<String>,
1351    options: Option<crate::QueryOptions>,
1352}
1353
1354impl Job for ResolveQueryJob {
1355    type Output = crate::QueryOptions;
1356    fn step(&mut self, response: Option<Response>) -> Result<Step<crate::QueryOptions>> {
1357        match self.inner.step(response)? {
1358            Step::Execute(r) => Ok(Step::Execute(r)),
1359            Step::Done(ticks) => {
1360                let mut o = self.options.take().expect("resolved once");
1361                for (r, t) in self.refs.iter().zip(&ticks) {
1362                    if o.as_of.as_deref() == Some(r.as_str()) {
1363                        o.as_of_tick = Some(*t);
1364                    }
1365                    o.versions.insert(r.clone(), *t);
1366                }
1367                Ok(Step::Done(o))
1368            }
1369        }
1370    }
1371}
1372
1373/// Changes the level of a store (one atomic request), then reloads its statistics.
1374pub fn level_change_job(
1375    state: &VersionState,
1376    to: Versioning,
1377    change: &LevelChange,
1378    caps: &Capabilities,
1379) -> Result<LevelChangeJob> {
1380    let statements = change_statements(state, to, change)?;
1381    Ok(LevelChangeJob {
1382        change: (!statements.is_empty()).then(|| Request::atomic(statements)),
1383        load: Some(crate::Stats::load_request(caps)),
1384    })
1385}
1386
1387pub struct LevelChangeJob {
1388    change: Option<Request>,
1389    load: Option<Request>,
1390}
1391
1392impl Job for LevelChangeJob {
1393    type Output = crate::Stats;
1394    fn step(&mut self, response: Option<Response>) -> Result<Step<crate::Stats>> {
1395        if let Some(r) = self.change.take() {
1396            return Ok(Step::Execute(r));
1397        }
1398        if let Some(r) = self.load.take() {
1399            return Ok(Step::Execute(r));
1400        }
1401        crate::Stats::from_response(&response.unwrap_or_default()).map(Step::Done)
1402    }
1403}
1404
1405// ---------------------------------------------------------------------------- history as data
1406
1407/// The vocabulary of the history graph (`GRAPH <oxilite:history>`).
1408pub mod vocab {
1409    /// `?commit oxl:added <<( s p o )>>`: a triple the commit added.
1410    pub const ADDED: &str = "https://oxilite.dev/ns#added";
1411    /// `?commit oxl:removed <<( s p o )>>`: a triple the commit removed.
1412    pub const REMOVED: &str = "https://oxilite.dev/ns#removed";
1413    pub const ACTIVITY: &str = "http://www.w3.org/ns/prov#Activity";
1414    pub const STARTED: &str = "http://www.w3.org/ns/prov#startedAtTime";
1415    pub const AGENT: &str = "http://www.w3.org/ns/prov#wasAssociatedWith";
1416    pub const INFORMED_BY: &str = "http://www.w3.org/ns/prov#wasInformedBy";
1417    pub const COMMENT: &str = "http://www.w3.org/2000/01/rdf-schema#comment";
1418    pub const TYPE: &str = "http://www.w3.org/1999/02/22-rdf-syntax-ns#type";
1419}
1420
1421/// SQL condition: the ticks shown as commits in the history (alias `k`): with a change log, the
1422/// ticks that changed something plus level changes; with only the clock, every tick.
1423fn activity_filter(state: &VersionState, k: &str) -> String {
1424    if state.history == History::None {
1425        "1".to_owned()
1426    } else {
1427        format!("({k}.kind <> 0 OR EXISTS (SELECT 1 FROM commits c WHERE c.tx = {k}.t))")
1428    }
1429}
1430
1431/// SQL: the `(s, p, o, g)` table of the history graph's commit descriptions. A commit is its
1432/// tick as an inline `xsd:integer`; its time, author and message are the terms its tick
1433/// recorded; `prov:wasInformedBy` links it to the commit before.
1434pub fn history_sql(state: &VersionState) -> Result<String> {
1435    if !state.has_ticks() {
1436        return Err(Error::Other(
1437            "this store keeps no history: raise its versioning level to `stamped` or `log`".into(),
1438        ));
1439    }
1440    let int = crate::compiler::expr_int_base();
1441    let id = crate::encoding::named_node_id;
1442    let f = activity_filter(state, "k");
1443    let fj = activity_filter(state, "j");
1444    let arm = |p: &str, o: &str, extra: &str| {
1445        format!(
1446            "SELECT k.t + {int} AS s, {} AS p, {o} AS o, 0 AS g FROM ticks k WHERE {f}{extra}",
1447            id(p)
1448        )
1449    };
1450    let parts = vec![
1451        arm(vocab::TYPE, &id(vocab::ACTIVITY).to_string(), ""),
1452        arm(vocab::STARTED, "k.time_id", " AND k.time_id IS NOT NULL"),
1453        arm(vocab::AGENT, "k.author_id", " AND k.author_id IS NOT NULL"),
1454        arm(
1455            vocab::COMMENT,
1456            "k.message_id",
1457            " AND k.message_id IS NOT NULL",
1458        ),
1459        arm(
1460            vocab::INFORMED_BY,
1461            &format!("(SELECT max(j.t) FROM ticks j WHERE j.t < k.t AND {fj}) + {int}"),
1462            &format!(" AND EXISTS (SELECT 1 FROM ticks j WHERE j.t < k.t AND {fj})"),
1463        ),
1464    ];
1465    Ok(format!("({})", crate::sql::union_all(parts, 5)))
1466}
1467
1468/// SQL: the Datalog relation `commit(c, parent, time, author)`: every commit (its tick as an
1469/// inline integer), the commit before it, and its time and author terms.
1470pub fn commits_rel_sql(state: &VersionState) -> String {
1471    let int = crate::compiler::expr_int_base();
1472    let (f, fj) = (activity_filter(state, "k"), activity_filter(state, "j"));
1473    format!(
1474        "(SELECT k.t + {int} AS c, (SELECT max(j.t) FROM ticks j WHERE j.t < k.t AND {fj}) + {int} AS parent, \
1475         k.time_id AS time, k.author_id AS author FROM ticks k WHERE {f})"
1476    )
1477}
1478
1479/// SQL: the Datalog relation `added(s, p, o, g, c)` or `removed(…)` from the change log.
1480pub fn changes_rel_sql(added: bool) -> String {
1481    let int = crate::compiler::expr_int_base();
1482    format!(
1483        "(SELECT s, p, o, g, tx + {int} AS c FROM quad_log WHERE op = {})",
1484        u8::from(added)
1485    )
1486}
1487
1488/// SQL: the Datalog relation `branch(name, c)`: `"main"` (term id `main`) and the latest tick.
1489pub fn branches_rel_sql(main: i64) -> String {
1490    let int = crate::compiler::expr_int_base();
1491    format!("(SELECT {main} AS name, max(t) + {int} AS c FROM ticks)")
1492}