Skip to main content

uqa_storage/sqlite/catalog/
cache_revisions.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Storage-owned, rollback-safe cache generations for every catalog session.
8
9use std::collections::BTreeMap;
10use std::fmt::Write as _;
11
12use super::{quote_sql_identifier, Catalog, Result, SQLiteError};
13use crate::CatalogCacheRevisions;
14
15impl Catalog {
16    pub(super) fn install_cache_revision_tracking(conn: &rusqlite::Connection) -> Result<()> {
17        let tables = conn
18            .prepare("SELECT name FROM sqlite_master WHERE type = 'table' AND substr(name, 1, 1) = '_' AND name NOT IN ('_cache_revisions', '_graph_path_pairs', '_graph_path_index_state') ORDER BY name")?
19            .query_map([], |row| row.get::<_, String>(0))?
20            .collect::<std::result::Result<Vec<_>, _>>()?;
21        for table in tables {
22            let columns = Self::table_columns(conn, &table)?.unwrap_or_default();
23            for (event, images) in [
24                ("INSERT", &["NEW"][..]),
25                ("DELETE", &["OLD"][..]),
26                ("UPDATE", &["OLD", "NEW"][..]),
27            ] {
28                let trigger = quote_sql_identifier(&format!("uqa_cache_{table}_{event}"));
29                let mut body = String::new();
30                for image in images {
31                    if matches!(table.as_str(), "_graph_vertices" | "_graph_edges") {
32                        let (entity_type, id) = if table == "_graph_vertices" {
33                            ("vertex", "vertex_id")
34                        } else {
35                            ("edge", "edge_id")
36                        };
37                        // A global entity can belong to several graphs. Its
38                        // property changes invalidate exactly those owners.
39                        write!(
40                            body,
41                            "INSERT INTO _cache_revisions(kind, name, generation) \
42                             SELECT 'graph', graph_name, 1 FROM _graph_membership \
43                             WHERE entity_type = '{entity_type}' AND entity_id = {image}.{id} \
44                             ON CONFLICT(kind, name) DO UPDATE SET generation = generation + 1;"
45                        )
46                        .expect("write graph revision trigger");
47                    } else {
48                        let (kind, name) =
49                            revision_scope(&table, columns.contains_key("table_name"), image);
50                        write!(body,
51                            "INSERT INTO _cache_revisions(kind, name, generation) VALUES ({kind}, {name}, 1) \
52                             ON CONFLICT(kind, name) DO UPDATE SET generation = generation + 1;"
53                        ).expect("writing a cache revision trigger to a String cannot fail");
54                    }
55                }
56                conn.execute_batch(&format!(
57                    "CREATE TRIGGER IF NOT EXISTS {trigger} AFTER {event} ON {} BEGIN {body} END;",
58                    quote_sql_identifier(&table),
59                ))?;
60            }
61        }
62        Ok(())
63    }
64
65    pub fn cache_revisions(&self) -> Result<CatalogCacheRevisions> {
66        self.conn.with(|conn| {
67            let mut revisions = CatalogCacheRevisions {
68                graphs: Some(BTreeMap::new()),
69                storage_schema: u64::from(conn.pragma_query_value(
70                    None,
71                    "schema_version",
72                    |row| row.get::<_, u32>(0),
73                )?),
74                ..CatalogCacheRevisions::default()
75            };
76            let mut statement = conn.prepare_cached(
77                "SELECT kind, name, generation FROM _cache_revisions ORDER BY kind, name",
78            )?;
79            let mut rows = statement.query([])?;
80            while let Some(row) = rows.next()? {
81                let kind: String = row.get(0)?;
82                let name: String = row.get(1)?;
83                let generation = u64::try_from(row.get::<_, i64>(2)?)
84                    .map_err(|_| SQLiteError::StorageBackend("negative cache revision".into()))?;
85                match kind.as_str() {
86                    "catalog" => revisions.table_catalog = generation,
87                    "registry" => revisions.registries = generation,
88                    "graph" => {
89                        revisions
90                            .graphs
91                            .as_mut()
92                            .expect("graph tracking enabled")
93                            .insert(name, generation);
94                    }
95                    "data" => {
96                        revisions.table_data.insert(name, generation);
97                    }
98                    "statistics" => {
99                        revisions.column_statistics.insert(name, generation);
100                    }
101                    "maintenance" => {
102                        revisions.statistics_maintenance.insert(name, generation);
103                    }
104                    _ => {
105                        return Err(SQLiteError::StorageBackend(format!(
106                            "unknown cache revision kind `{kind}`"
107                        )))
108                    }
109                }
110            }
111            Ok(revisions)
112        })
113    }
114}
115
116fn revision_scope(table: &str, has_table_name: bool, image: &str) -> (String, String) {
117    match table {
118        "_tables" => ("'catalog'".into(), "''".into()),
119        "_column_stats" => ("'statistics'".into(), format!("{image}.table_name")),
120        "_named_graphs" => ("'graph'".into(), format!("{image}.name")),
121        "_graph_membership" => ("'graph'".into(), format!("{image}.graph_name")),
122        "_metadata" => {
123            let key = format!("{image}.key");
124            let maintenance = "uqa.statistics.maintenance.v1:";
125            let next_id = "uqa.table_next_id.v1:";
126            let graph_labels = "graph_label_registry::";
127            (
128                format!("CASE WHEN substr({key}, 1, {}) = '{maintenance}' THEN 'maintenance' WHEN substr({key}, 1, {}) = '{next_id}' THEN 'data' WHEN substr({key}, 1, {}) = '{graph_labels}' THEN 'graph' ELSE 'registry' END", maintenance.len(), next_id.len(), graph_labels.len()),
129                format!("CASE WHEN substr({key}, 1, {}) = '{maintenance}' THEN substr({key}, {}) WHEN substr({key}, 1, {}) = '{next_id}' THEN substr({key}, {}) WHEN substr({key}, 1, {}) = '{graph_labels}' THEN substr({key}, {}) ELSE '' END", maintenance.len(), maintenance.len() + 1, next_id.len(), next_id.len() + 1, graph_labels.len(), graph_labels.len() + 1),
130            )
131        }
132        // These rows define access paths, not the data stored in those paths.
133        "_table_field_analyzers" | "_catalog_indexes" | "_btree_indexes" => {
134            ("'registry'".into(), "''".into())
135        }
136        _ if has_table_name => ("'data'".into(), format!("{image}.table_name")),
137        _ => ("'registry'".into(), "''".into()),
138    }
139}