Skip to main content

uqa_storage/sqlite/catalog/
schema_tables.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Metadata, schema, table, and column lifecycle.
8
9use super::{
10    columns_json_references, delete_table_rows_if_exists, drop_fts_aux_tables_for_field,
11    drop_fts_aux_tables_for_table, migration_relation, params,
12    rename_btree_field_rows_or_keep_existing, rename_field_rows_or_keep_existing,
13    rename_fts_aux_tables_for_field, renamed_columns_json, table_exists,
14    update_btree_table_name_rows_if_exists, update_table_name_rows_if_exists, Catalog,
15    OptionalExtension, RelationIdentity, RelationKind, Result, SQLiteError, SchemaRow, TableSchema,
16    VectorFieldSchema,
17};
18
19impl Catalog {
20    /// Store an arbitrary key/value pair in the `_metadata` table.
21    pub fn set_metadata(&self, key: &str, value: &str) -> Result<()> {
22        self.conn.with(|c| {
23            c.execute(
24                "INSERT OR REPLACE INTO _metadata (key, value) VALUES (?1, ?2)",
25                params![key, value],
26            )?;
27            Ok(())
28        })
29    }
30
31    /// Read a key/value pair from the `_metadata` table.
32    pub fn get_metadata(&self, key: &str) -> Result<Option<String>> {
33        self.conn.with(|c| {
34            let v: Option<String> = c
35                .query_row(
36                    "SELECT value FROM _metadata WHERE key = ?1",
37                    params![key],
38                    |r| r.get(0),
39                )
40                .optional()?;
41            Ok(v)
42        })
43    }
44
45    pub fn save_schema(&self, name: &str) -> Result<()> {
46        self.save_schema_row(&SchemaRow::legacy(name))
47    }
48
49    pub fn save_schema_row(&self, schema: &SchemaRow) -> Result<()> {
50        let acl_json = schema.acl.as_ref().map(serde_json::to_string).transpose()?;
51        self.conn.with(|c| {
52            c.execute(
53                "INSERT INTO _schemas (name, role_owner, acl_json) VALUES (?1, ?2, ?3)
54                 ON CONFLICT(name) DO UPDATE SET role_owner = excluded.role_owner, acl_json = excluded.acl_json",
55                params![schema.name, schema.role_owner, acl_json],
56            )?;
57            Ok(())
58        })
59    }
60
61    pub fn drop_schema(&self, name: &str) -> Result<()> {
62        self.conn.with(|c| {
63            let relation_count: i64 = c.query_row(
64                "SELECT COUNT(*) FROM _relations WHERE schema_name = ?1",
65                params![name],
66                |row| row.get(0),
67            )?;
68            if relation_count != 0 {
69                return Err(SQLiteError::StorageBackend(format!(
70                    "schema `{name}` still owns catalog relations"
71                )));
72            }
73            c.execute("DELETE FROM _schemas WHERE name = ?1", params![name])?;
74            Ok(())
75        })
76    }
77
78    pub fn load_schemas(&self) -> Result<Vec<String>> {
79        Ok(self
80            .load_schema_rows()?
81            .into_iter()
82            .map(|schema| schema.name)
83            .collect())
84    }
85
86    pub fn load_schema_rows(&self) -> Result<Vec<SchemaRow>> {
87        self.conn.with(|c| {
88            let mut stmt =
89                c.prepare("SELECT name, role_owner, acl_json FROM _schemas ORDER BY name")?;
90            let rows = stmt.query_map([], |row| {
91                Ok((
92                    row.get::<_, String>(0)?,
93                    row.get::<_, String>(1)?,
94                    row.get::<_, Option<String>>(2)?,
95                ))
96            })?;
97            let mut out = Vec::new();
98            for row in rows {
99                let (name, role_owner, acl_json) = row?;
100                let acl = acl_json
101                    .map(|json| serde_json::from_str(&json))
102                    .transpose()?;
103                out.push(SchemaRow {
104                    name,
105                    role_owner,
106                    acl,
107                });
108            }
109            Ok(out)
110        })
111    }
112
113    pub fn save_table(&self, schema: &TableSchema) -> Result<()> {
114        let analyzer = schema.analyzer_json.clone();
115        let fts = serde_json::to_string(&schema.fts_fields)?;
116        let vectors = serde_json::to_string(&schema.vector_fields)?;
117        let columns = schema.columns_json.clone();
118        let constraints = schema.constraints_json.clone();
119        let role_owner = schema.role_owner.clone();
120        let acl_json = schema.acl.as_ref().map(serde_json::to_string).transpose()?;
121        let column_acls_json = serde_json::to_string(&schema.column_acls)?;
122        let object_id = schema.object_id;
123        let storage_generation = schema.storage_generation;
124        self.conn.with_mut(|c| {
125            let tx = c.savepoint()?;
126            Self::claim_relation(&tx, &schema.relation, RelationKind::Table)?;
127            tx.execute(
128                "INSERT INTO _tables
129                    (schema_name, relation_name, kind, analyzer, fts_fields,
130                     vector_fields, columns, constraints, storage_generation, object_id,
131                     role_owner, acl_json, column_acls_json)
132                 VALUES (?1, ?2, 'table', ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)
133                 ON CONFLICT(schema_name, relation_name) DO UPDATE SET
134                     analyzer = excluded.analyzer,
135                     fts_fields = excluded.fts_fields,
136                     vector_fields = excluded.vector_fields,
137                     columns = excluded.columns,
138                     constraints = excluded.constraints,
139                     storage_generation = excluded.storage_generation,
140                     object_id = excluded.object_id,
141                     role_owner = excluded.role_owner,
142                     acl_json = excluded.acl_json,
143                     column_acls_json = excluded.column_acls_json",
144                params![
145                    schema.relation.schema,
146                    schema.relation.name,
147                    analyzer,
148                    fts,
149                    vectors,
150                    columns,
151                    constraints,
152                    storage_generation.as_slice(),
153                    object_id.as_slice(),
154                    role_owner,
155                    acl_json,
156                    column_acls_json
157                ],
158            )?;
159            tx.commit()?;
160            Ok(())
161        })
162    }
163
164    pub fn load_tables(&self) -> Result<Vec<TableSchema>> {
165        self.conn.with(|c| {
166            let mut stmt = c.prepare(
167                "SELECT schema_name, relation_name, analyzer, fts_fields,
168                        vector_fields, columns, constraints, storage_generation, object_id,
169                        role_owner, acl_json, column_acls_json
170                   FROM _tables ORDER BY schema_name, relation_name",
171            )?;
172            let rows = stmt.query_map([], |r| {
173                Ok((
174                    r.get::<_, String>(0)?,
175                    r.get::<_, String>(1)?,
176                    r.get::<_, String>(2)?,
177                    r.get::<_, String>(3)?,
178                    r.get::<_, String>(4)?,
179                    r.get::<_, Option<String>>(5)?,
180                    r.get::<_, String>(6)?,
181                    r.get::<_, Vec<u8>>(7)?,
182                    r.get::<_, Vec<u8>>(8)?,
183                    r.get::<_, String>(9)?,
184                    r.get::<_, Option<String>>(10)?,
185                    r.get::<_, Option<String>>(11)?,
186                ))
187            })?;
188            let mut out = Vec::new();
189            for row in rows {
190                let (
191                    schema_name,
192                    relation_name,
193                    analyzer_json,
194                    fts_str,
195                    vec_str,
196                    cols_opt,
197                    constraints_json,
198                    storage_generation,
199                    object_id,
200                    role_owner,
201                    acl_json,
202                    column_acls_json,
203                ) = row?;
204                let fts_fields: Vec<String> = serde_json::from_str(&fts_str)?;
205                let vector_fields: Vec<VectorFieldSchema> = serde_json::from_str(&vec_str)?;
206                let storage_generation: [u8; 16] = storage_generation.try_into().map_err(|value: Vec<u8>| {
207                    SQLiteError::StorageBackend(format!(
208                        "table `{schema_name}.{relation_name}` has a {}-byte storage generation instead of 16 bytes",
209                        value.len()
210                    ))
211                })?;
212                let object_id: [u8; 16] = object_id.try_into().map_err(|value: Vec<u8>| {
213                    SQLiteError::StorageBackend(format!(
214                        "table `{schema_name}.{relation_name}` has a {}-byte object identity instead of 16 bytes",
215                        value.len()
216                    ))
217                })?;
218                let acl = acl_json
219                    .map(|json| serde_json::from_str(&json))
220                    .transpose()?;
221                let column_acls = column_acls_json
222                    .map(|json| serde_json::from_str(&json))
223                    .transpose()?
224                    .unwrap_or_default();
225                out.push(TableSchema {
226                    relation: RelationIdentity::new(schema_name, relation_name),
227                    role_owner,
228                    acl,
229                    column_acls,
230                    object_id,
231                    storage_generation,
232                    analyzer_json,
233                    fts_fields,
234                    vector_fields,
235                    columns_json: cols_opt.unwrap_or_default(),
236                    constraints_json,
237                });
238            }
239            Ok(out)
240        })
241    }
242
243    pub fn drop_table(&self, name: &str) -> Result<()> {
244        let relation = migration_relation(name)?;
245        self.conn.with_mut(|c| {
246            let tx = c.savepoint()?;
247            Self::drop_catalog_index_rows_for_table(&tx, &relation)?;
248            tx.execute(
249                "DELETE FROM _tables WHERE schema_name = ?1 AND relation_name = ?2",
250                params![relation.schema, relation.name],
251            )?;
252            Self::release_relation(&tx, &relation, RelationKind::Table)?;
253            tx.commit()?;
254            Ok(())
255        })
256    }
257
258    /// Wipe the rows owned by `table` from the per-table data tables
259    /// (`_documents`, clustered postings, `_doc_lengths`, `_field_stats`,
260    /// `_vectors`, IVF/HNSW metadata). Run after [`Catalog::drop_table`]
261    /// when the engine drops the table from its in-memory registry as
262    /// well.
263    pub fn purge_table_data(&self, name: &str) -> Result<()> {
264        let relation = migration_relation(name)?;
265        let storage_names = relation.canonical_and_legacy_public_names();
266        self.conn.with_mut(|c| {
267            let tx = c.savepoint()?;
268            for storage_name in &storage_names {
269                for table in [
270                    "_documents",
271                    "_document_blobs",
272                    "_posting_clusters",
273                    "_posting_documents",
274                    "_doc_lengths",
275                    "_field_stats",
276                    "_vectors",
277                    "_ivf_indexes",
278                    "_ivf_centroids",
279                    "_ivf_assignments",
280                    "_hnsw_indexes",
281                    "_hnsw_nodes",
282                    "_hnsw_edges",
283                    "_column_stats",
284                    "_btree_index_entries",
285                    "_btree_indexes",
286                ] {
287                    delete_table_rows_if_exists(&tx, table, storage_name)?;
288                }
289                drop_fts_aux_tables_for_table(&tx, storage_name)?;
290            }
291            tx.commit()?;
292            Ok(())
293        })
294    }
295
296    pub fn drop_table_and_data(&self, name: &str) -> Result<()> {
297        let relation = migration_relation(name)?;
298        let storage_names = relation.canonical_and_legacy_public_names();
299        self.conn.with_mut(|c| {
300            let tx = c.savepoint()?;
301            Self::drop_catalog_index_rows_for_table(&tx, &relation)?;
302            tx.execute(
303                "DELETE FROM _tables WHERE schema_name = ?1 AND relation_name = ?2",
304                params![relation.schema, relation.name],
305            )?;
306            for storage_name in &storage_names {
307                for table in [
308                    "_documents",
309                    "_document_blobs",
310                    "_posting_clusters",
311                    "_posting_documents",
312                    "_doc_lengths",
313                    "_field_stats",
314                    "_vectors",
315                    "_ivf_indexes",
316                    "_ivf_centroids",
317                    "_ivf_assignments",
318                    "_hnsw_indexes",
319                    "_hnsw_nodes",
320                    "_hnsw_edges",
321                    "_column_stats",
322                    "_btree_index_entries",
323                    "_btree_indexes",
324                ] {
325                    delete_table_rows_if_exists(&tx, table, storage_name)?;
326                }
327                tx.execute(
328                    "DELETE FROM _table_field_analyzers WHERE table_name = ?1",
329                    params![storage_name],
330                )?;
331                drop_fts_aux_tables_for_table(&tx, storage_name)?;
332            }
333            Self::release_relation(&tx, &relation, RelationKind::Table)?;
334            tx.commit()?;
335            Ok(())
336        })
337    }
338
339    pub fn rename_table_data(&self, from: &str, to: &str) -> Result<()> {
340        let from_relation = migration_relation(from)?;
341        let to_relation = migration_relation(to)?;
342        if from_relation == to_relation {
343            return Ok(());
344        }
345        if from_relation.schema != to_relation.schema {
346            return Err(SQLiteError::StorageBackend(
347                "moving a table between schemas is not supported by the catalog".into(),
348            ));
349        }
350        self.conn.with_mut(|c| {
351            let tx = c.savepoint()?;
352            Self::claim_relation(&tx, &to_relation, RelationKind::Table)?;
353            let updated = tx.execute(
354                "UPDATE _tables
355                    SET schema_name = ?3, relation_name = ?4
356                  WHERE schema_name = ?1 AND relation_name = ?2",
357                params![
358                    from_relation.schema,
359                    from_relation.name,
360                    to_relation.schema,
361                    to_relation.name
362                ],
363            )?;
364            if updated == 0 {
365                return Err(SQLiteError::StorageBackend(format!(
366                    "table `{from}` does not exist"
367                )));
368            }
369            for table in [
370                "_documents",
371                "_document_blobs",
372                "_posting_clusters",
373                "_posting_documents",
374                "_doc_lengths",
375                "_field_stats",
376                "_vectors",
377                "_ivf_indexes",
378                "_ivf_centroids",
379                "_ivf_assignments",
380                "_hnsw_indexes",
381                "_hnsw_nodes",
382                "_hnsw_edges",
383                "_column_stats",
384                "_table_field_analyzers",
385            ] {
386                update_table_name_rows_if_exists(&tx, table, from, to)?;
387            }
388            update_btree_table_name_rows_if_exists(&tx, from, to)?;
389            drop_fts_aux_tables_for_table(&tx, from)?;
390            Self::release_relation(&tx, &from_relation, RelationKind::Table)?;
391            tx.commit()?;
392            Ok(())
393        })
394    }
395
396    pub fn drop_column_data(&self, table_name: &str, column_name: &str) -> Result<()> {
397        let indexes = self.catalog_indexes_referencing_column(table_name, column_name)?;
398        self.conn.with_mut(|c| {
399            let tx = c.savepoint()?;
400            if table_exists(&tx, "_document_blobs")? {
401                tx.execute(
402                    "DELETE FROM _document_blobs WHERE table_name = ?1 AND field_name = ?2",
403                    params![table_name, column_name],
404                )?;
405            }
406            for table in [
407                "_posting_clusters",
408                "_posting_documents",
409                "_doc_lengths",
410                "_field_stats",
411                "_vectors",
412                "_ivf_indexes",
413                "_ivf_centroids",
414                "_ivf_assignments",
415                "_hnsw_indexes",
416                "_hnsw_nodes",
417                "_hnsw_edges",
418                "_btree_index_entries",
419                "_btree_indexes",
420            ] {
421                tx.execute(
422                    &format!("DELETE FROM {table} WHERE table_name = ?1 AND field = ?2"),
423                    params![table_name, column_name],
424                )?;
425            }
426            tx.execute(
427                "DELETE FROM _column_stats WHERE table_name = ?1 AND column_name = ?2",
428                params![table_name, column_name],
429            )?;
430            tx.execute(
431                "DELETE FROM _table_field_analyzers WHERE table_name = ?1 AND field = ?2",
432                params![table_name, column_name],
433            )?;
434            for index in indexes {
435                tx.execute(
436                    "DELETE FROM _catalog_indexes
437                      WHERE schema_name = ?1 AND relation_name = ?2",
438                    params![index.schema, index.name],
439                )?;
440                Self::release_relation(&tx, &index, RelationKind::Index)?;
441            }
442            drop_fts_aux_tables_for_field(&tx, table_name, column_name)?;
443            tx.commit()?;
444            Ok(())
445        })
446    }
447
448    pub fn rename_column_data(&self, table_name: &str, from: &str, to: &str) -> Result<()> {
449        let index_updates = self.catalog_index_column_renames(table_name, from, to)?;
450        self.conn.with_mut(|c| {
451            let tx = c.savepoint()?;
452            rename_field_rows_or_keep_existing(
453                &tx,
454                "_document_blobs",
455                "field_name",
456                table_name,
457                from,
458                to,
459            )?;
460            for table in [
461                "_posting_clusters",
462                "_posting_documents",
463                "_doc_lengths",
464                "_field_stats",
465                "_vectors",
466                "_ivf_indexes",
467                "_ivf_centroids",
468                "_ivf_assignments",
469                "_hnsw_indexes",
470                "_hnsw_nodes",
471                "_hnsw_edges",
472            ] {
473                rename_field_rows_or_keep_existing(&tx, table, "field", table_name, from, to)?;
474            }
475            rename_btree_field_rows_or_keep_existing(&tx, table_name, from, to)?;
476            rename_field_rows_or_keep_existing(
477                &tx,
478                "_column_stats",
479                "column_name",
480                table_name,
481                from,
482                to,
483            )?;
484            rename_field_rows_or_keep_existing(
485                &tx,
486                "_table_field_analyzers",
487                "field",
488                table_name,
489                from,
490                to,
491            )?;
492            for (index, columns_json) in index_updates {
493                tx.execute(
494                    "UPDATE _catalog_indexes
495                        SET columns = ?2
496                      WHERE schema_name = ?1 AND relation_name = ?3",
497                    params![index.schema, columns_json, index.name],
498                )?;
499            }
500            rename_fts_aux_tables_for_field(&tx, table_name, from, to)?;
501            tx.commit()?;
502            Ok(())
503        })
504    }
505
506    pub(super) fn catalog_indexes_referencing_column(
507        &self,
508        table_name: &str,
509        column_name: &str,
510    ) -> Result<Vec<RelationIdentity>> {
511        let mut out = Vec::new();
512        for row in self.load_catalog_indexes()? {
513            if row.table_name == table_name
514                && columns_json_references(&row.columns_json, column_name)?
515            {
516                out.push(row.relation);
517            }
518        }
519        Ok(out)
520    }
521
522    pub(super) fn catalog_index_column_renames(
523        &self,
524        table_name: &str,
525        from: &str,
526        to: &str,
527    ) -> Result<Vec<(RelationIdentity, String)>> {
528        let mut out = Vec::new();
529        for row in self.load_catalog_indexes()? {
530            if row.table_name != table_name {
531                continue;
532            }
533            if let Some(columns_json) = renamed_columns_json(&row.columns_json, from, to)? {
534                out.push((row.relation, columns_json));
535            }
536        }
537        Ok(out)
538    }
539}