1use 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#[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 #[default]
32 Off,
33 Stamped,
35 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#[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 #[default]
86 None,
87 Live,
89 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#[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 pub stamp_column: bool,
113 pub stamp_index: bool,
115 pub as_of_index: bool,
117}
118
119impl VersionState {
120 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 pub fn has_ticks(&self) -> bool {
140 self.level >= Versioning::Stamped || self.stamp_column || self.history != History::None
141 }
142}
143
144#[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#[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 pub as_of_index: Option<bool>,
160 pub stamp_index: Option<bool>,
162 pub allow_loss: bool,
164 pub author: Option<String>,
166 pub message: Option<String>,
167}
168
169pub mod kind {
172 pub const WRITE: i64 = 0;
174 pub const GENESIS: i64 = 1;
176 pub const FREEZE: i64 = 2;
178 pub const RESUME: i64 = 3;
180 pub const DROPPED: i64 = 4;
182 pub const LEVEL: i64 = 5;
184 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
201pub const CURRENT_TICK: &str = "(SELECT max(t) FROM ticks)";
203
204pub const NOW: &str = "((julianday('now') - 2440587.5) * 86400.0)";
206
207pub 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
216pub 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
229pub const PREFIX_LEN: usize = 2;
232
233pub 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
266pub 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
286pub fn strip(mut response: Response) -> Response {
288 response.drain(..PREFIX_LEN.min(response.len()));
289 response
290}
291
292pub 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
303pub 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
344pub 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 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
446fn 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
461fn clock_ddl() -> Vec<Statement> {
463 [
464 "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 "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
476fn 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 "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
494fn 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 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
546pub 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 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
702pub 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#[derive(Debug, Clone, PartialEq)]
715pub enum VersionRef {
716 Head,
718 Back(u64),
720 Tick(i64),
722 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
786pub 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
865pub 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 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
888fn 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#[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 pub head: Option<i64>,
915 pub head_time: Option<f64>,
917 pub genesis: Option<i64>,
919 pub frozen_at: Option<i64>,
921 pub commits: Option<u64>,
923}
924
925pub 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#[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 pub time: f64,
974 pub kind: String,
976 pub author: Option<String>,
977 pub message: Option<String>,
978 pub added: Option<u64>,
980 pub removed: Option<u64>,
981}
982
983pub 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#[derive(Debug, Clone, PartialEq)]
1032pub struct Change {
1033 pub tick: i64,
1034 pub added: bool,
1035 pub quad: oxrdf::Quad,
1036}
1037
1038pub 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
1117pub 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
1159pub 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
1191pub 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
1227pub 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
1252pub fn query_version_refs(query: &spargebra::Query, options: &crate::QueryOptions) -> Vec<String> {
1255 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
1275pub 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
1331pub 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
1373pub 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
1405pub mod vocab {
1409 pub const ADDED: &str = "https://oxilite.dev/ns#added";
1411 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
1421fn 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
1431pub 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
1468pub 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
1479pub 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
1488pub 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}