use super::{
columns_json_references, delete_table_rows_if_exists, drop_fts_aux_tables_for_field,
drop_fts_aux_tables_for_table, migration_relation, params,
rename_btree_field_rows_or_keep_existing, rename_field_rows_or_keep_existing,
rename_fts_aux_tables_for_field, renamed_columns_json, table_exists,
update_btree_table_name_rows_if_exists, update_table_name_rows_if_exists, Catalog,
OptionalExtension, RelationIdentity, RelationKind, Result, SQLiteError, SchemaRow, TableSchema,
VectorFieldSchema,
};
impl Catalog {
pub fn set_metadata(&self, key: &str, value: &str) -> Result<()> {
self.conn.with(|c| {
c.execute(
"INSERT OR REPLACE INTO _metadata (key, value) VALUES (?1, ?2)",
params![key, value],
)?;
Ok(())
})
}
pub fn get_metadata(&self, key: &str) -> Result<Option<String>> {
self.conn.with(|c| {
let v: Option<String> = c
.query_row(
"SELECT value FROM _metadata WHERE key = ?1",
params![key],
|r| r.get(0),
)
.optional()?;
Ok(v)
})
}
pub fn save_schema(&self, name: &str) -> Result<()> {
self.save_schema_row(&SchemaRow::legacy(name))
}
pub fn save_schema_row(&self, schema: &SchemaRow) -> Result<()> {
let acl_json = schema.acl.as_ref().map(serde_json::to_string).transpose()?;
self.conn.with(|c| {
c.execute(
"INSERT INTO _schemas (name, role_owner, acl_json) VALUES (?1, ?2, ?3)
ON CONFLICT(name) DO UPDATE SET role_owner = excluded.role_owner, acl_json = excluded.acl_json",
params![schema.name, schema.role_owner, acl_json],
)?;
Ok(())
})
}
pub fn drop_schema(&self, name: &str) -> Result<()> {
self.conn.with(|c| {
let relation_count: i64 = c.query_row(
"SELECT COUNT(*) FROM _relations WHERE schema_name = ?1",
params![name],
|row| row.get(0),
)?;
if relation_count != 0 {
return Err(SQLiteError::StorageBackend(format!(
"schema `{name}` still owns catalog relations"
)));
}
c.execute("DELETE FROM _schemas WHERE name = ?1", params![name])?;
Ok(())
})
}
pub fn load_schemas(&self) -> Result<Vec<String>> {
Ok(self
.load_schema_rows()?
.into_iter()
.map(|schema| schema.name)
.collect())
}
pub fn load_schema_rows(&self) -> Result<Vec<SchemaRow>> {
self.conn.with(|c| {
let mut stmt =
c.prepare("SELECT name, role_owner, acl_json FROM _schemas ORDER BY name")?;
let rows = stmt.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, Option<String>>(2)?,
))
})?;
let mut out = Vec::new();
for row in rows {
let (name, role_owner, acl_json) = row?;
let acl = acl_json
.map(|json| serde_json::from_str(&json))
.transpose()?;
out.push(SchemaRow {
name,
role_owner,
acl,
});
}
Ok(out)
})
}
pub fn save_table(&self, schema: &TableSchema) -> Result<()> {
let analyzer = schema.analyzer_json.clone();
let fts = serde_json::to_string(&schema.fts_fields)?;
let vectors = serde_json::to_string(&schema.vector_fields)?;
let columns = schema.columns_json.clone();
let constraints = schema.constraints_json.clone();
let role_owner = schema.role_owner.clone();
let acl_json = schema.acl.as_ref().map(serde_json::to_string).transpose()?;
let column_acls_json = serde_json::to_string(&schema.column_acls)?;
let object_id = schema.object_id;
let storage_generation = schema.storage_generation;
self.conn.with_mut(|c| {
let tx = c.savepoint()?;
Self::claim_relation(&tx, &schema.relation, RelationKind::Table)?;
tx.execute(
"INSERT INTO _tables
(schema_name, relation_name, kind, analyzer, fts_fields,
vector_fields, columns, constraints, storage_generation, object_id,
role_owner, acl_json, column_acls_json)
VALUES (?1, ?2, 'table', ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)
ON CONFLICT(schema_name, relation_name) DO UPDATE SET
analyzer = excluded.analyzer,
fts_fields = excluded.fts_fields,
vector_fields = excluded.vector_fields,
columns = excluded.columns,
constraints = excluded.constraints,
storage_generation = excluded.storage_generation,
object_id = excluded.object_id,
role_owner = excluded.role_owner,
acl_json = excluded.acl_json,
column_acls_json = excluded.column_acls_json",
params![
schema.relation.schema,
schema.relation.name,
analyzer,
fts,
vectors,
columns,
constraints,
storage_generation.as_slice(),
object_id.as_slice(),
role_owner,
acl_json,
column_acls_json
],
)?;
tx.commit()?;
Ok(())
})
}
pub fn load_tables(&self) -> Result<Vec<TableSchema>> {
self.conn.with(|c| {
let mut stmt = c.prepare(
"SELECT schema_name, relation_name, analyzer, fts_fields,
vector_fields, columns, constraints, storage_generation, object_id,
role_owner, acl_json, column_acls_json
FROM _tables ORDER BY schema_name, relation_name",
)?;
let rows = stmt.query_map([], |r| {
Ok((
r.get::<_, String>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
r.get::<_, String>(3)?,
r.get::<_, String>(4)?,
r.get::<_, Option<String>>(5)?,
r.get::<_, String>(6)?,
r.get::<_, Vec<u8>>(7)?,
r.get::<_, Vec<u8>>(8)?,
r.get::<_, String>(9)?,
r.get::<_, Option<String>>(10)?,
r.get::<_, Option<String>>(11)?,
))
})?;
let mut out = Vec::new();
for row in rows {
let (
schema_name,
relation_name,
analyzer_json,
fts_str,
vec_str,
cols_opt,
constraints_json,
storage_generation,
object_id,
role_owner,
acl_json,
column_acls_json,
) = row?;
let fts_fields: Vec<String> = serde_json::from_str(&fts_str)?;
let vector_fields: Vec<VectorFieldSchema> = serde_json::from_str(&vec_str)?;
let storage_generation: [u8; 16] = storage_generation.try_into().map_err(|value: Vec<u8>| {
SQLiteError::StorageBackend(format!(
"table `{schema_name}.{relation_name}` has a {}-byte storage generation instead of 16 bytes",
value.len()
))
})?;
let object_id: [u8; 16] = object_id.try_into().map_err(|value: Vec<u8>| {
SQLiteError::StorageBackend(format!(
"table `{schema_name}.{relation_name}` has a {}-byte object identity instead of 16 bytes",
value.len()
))
})?;
let acl = acl_json
.map(|json| serde_json::from_str(&json))
.transpose()?;
let column_acls = column_acls_json
.map(|json| serde_json::from_str(&json))
.transpose()?
.unwrap_or_default();
out.push(TableSchema {
relation: RelationIdentity::new(schema_name, relation_name),
role_owner,
acl,
column_acls,
object_id,
storage_generation,
analyzer_json,
fts_fields,
vector_fields,
columns_json: cols_opt.unwrap_or_default(),
constraints_json,
});
}
Ok(out)
})
}
pub fn drop_table(&self, name: &str) -> Result<()> {
let relation = migration_relation(name)?;
self.conn.with_mut(|c| {
let tx = c.savepoint()?;
Self::drop_catalog_index_rows_for_table(&tx, &relation)?;
tx.execute(
"DELETE FROM _tables WHERE schema_name = ?1 AND relation_name = ?2",
params![relation.schema, relation.name],
)?;
Self::release_relation(&tx, &relation, RelationKind::Table)?;
tx.commit()?;
Ok(())
})
}
pub fn purge_table_data(&self, name: &str) -> Result<()> {
let relation = migration_relation(name)?;
let storage_names = relation.canonical_and_legacy_public_names();
self.conn.with_mut(|c| {
let tx = c.savepoint()?;
for storage_name in &storage_names {
for table in [
"_documents",
"_document_blobs",
"_postings",
"_posting_clusters",
"_posting_documents",
"_doc_lengths",
"_field_stats",
"_occurrence_clusters",
"_occurrence_documents",
"_occurrence_lengths",
"_occurrence_fields",
"_occurrence_formats",
"_vectors",
"_ivf_indexes",
"_ivf_centroids",
"_ivf_assignments",
"_hnsw_indexes",
"_hnsw_nodes",
"_hnsw_edges",
"_column_stats",
"_btree_index_entries",
"_btree_indexes",
] {
delete_table_rows_if_exists(&tx, table, storage_name)?;
}
drop_fts_aux_tables_for_table(&tx, storage_name)?;
}
tx.commit()?;
Ok(())
})
}
pub fn drop_table_and_data(&self, name: &str) -> Result<()> {
let relation = migration_relation(name)?;
let storage_names = relation.canonical_and_legacy_public_names();
self.conn.with_mut(|c| {
let tx = c.savepoint()?;
Self::drop_catalog_index_rows_for_table(&tx, &relation)?;
tx.execute(
"DELETE FROM _tables WHERE schema_name = ?1 AND relation_name = ?2",
params![relation.schema, relation.name],
)?;
for storage_name in &storage_names {
for table in [
"_documents",
"_document_blobs",
"_postings",
"_posting_clusters",
"_posting_documents",
"_doc_lengths",
"_field_stats",
"_occurrence_clusters",
"_occurrence_documents",
"_occurrence_lengths",
"_occurrence_fields",
"_occurrence_formats",
"_vectors",
"_ivf_indexes",
"_ivf_centroids",
"_ivf_assignments",
"_hnsw_indexes",
"_hnsw_nodes",
"_hnsw_edges",
"_column_stats",
"_btree_index_entries",
"_btree_indexes",
] {
delete_table_rows_if_exists(&tx, table, storage_name)?;
}
tx.execute(
"DELETE FROM _table_field_analyzers WHERE table_name = ?1",
params![storage_name],
)?;
drop_fts_aux_tables_for_table(&tx, storage_name)?;
}
Self::release_relation(&tx, &relation, RelationKind::Table)?;
tx.commit()?;
Ok(())
})
}
pub fn rename_table_data(&self, from: &str, to: &str) -> Result<()> {
let from_relation = migration_relation(from)?;
let to_relation = migration_relation(to)?;
if from_relation == to_relation {
return Ok(());
}
if from_relation.schema != to_relation.schema {
return Err(SQLiteError::StorageBackend(
"moving a table between schemas is not supported by the catalog".into(),
));
}
self.conn.with_mut(|c| {
let tx = c.savepoint()?;
Self::claim_relation(&tx, &to_relation, RelationKind::Table)?;
let updated = tx.execute(
"UPDATE _tables
SET schema_name = ?3, relation_name = ?4
WHERE schema_name = ?1 AND relation_name = ?2",
params![
from_relation.schema,
from_relation.name,
to_relation.schema,
to_relation.name
],
)?;
if updated == 0 {
return Err(SQLiteError::StorageBackend(format!(
"table `{from}` does not exist"
)));
}
for table in [
"_documents",
"_document_blobs",
"_postings",
"_posting_clusters",
"_posting_documents",
"_doc_lengths",
"_field_stats",
"_occurrence_clusters",
"_occurrence_documents",
"_occurrence_lengths",
"_occurrence_fields",
"_occurrence_formats",
"_vectors",
"_ivf_indexes",
"_ivf_centroids",
"_ivf_assignments",
"_hnsw_indexes",
"_hnsw_nodes",
"_hnsw_edges",
"_column_stats",
"_table_field_analyzers",
] {
update_table_name_rows_if_exists(&tx, table, from, to)?;
}
update_btree_table_name_rows_if_exists(&tx, from, to)?;
drop_fts_aux_tables_for_table(&tx, from)?;
Self::release_relation(&tx, &from_relation, RelationKind::Table)?;
tx.commit()?;
Ok(())
})
}
pub fn drop_column_data(&self, table_name: &str, column_name: &str) -> Result<()> {
let indexes = self.catalog_indexes_referencing_column(table_name, column_name)?;
self.conn.with_mut(|c| {
let tx = c.savepoint()?;
if table_exists(&tx, "_document_blobs")? {
tx.execute(
"DELETE FROM _document_blobs WHERE table_name = ?1 AND field_name = ?2",
params![table_name, column_name],
)?;
}
for table in [
"_postings",
"_posting_clusters",
"_posting_documents",
"_doc_lengths",
"_field_stats",
"_occurrence_clusters",
"_occurrence_documents",
"_occurrence_lengths",
"_occurrence_fields",
"_vectors",
"_ivf_indexes",
"_ivf_centroids",
"_ivf_assignments",
"_hnsw_indexes",
"_hnsw_nodes",
"_hnsw_edges",
"_btree_index_entries",
"_btree_indexes",
] {
if matches!(
table,
"_postings"
| "_posting_clusters"
| "_posting_documents"
| "_doc_lengths"
| "_field_stats"
) && !Self::table_columns(&tx, table)?
.is_some_and(|columns| columns.contains_key("field"))
{
continue;
}
tx.execute(
&format!("DELETE FROM {table} WHERE table_name = ?1 AND field = ?2"),
params![table_name, column_name],
)?;
}
tx.execute(
"DELETE FROM _column_stats WHERE table_name = ?1 AND column_name = ?2",
params![table_name, column_name],
)?;
tx.execute(
"DELETE FROM _table_field_analyzers WHERE table_name = ?1 AND field = ?2",
params![table_name, column_name],
)?;
for index in indexes {
tx.execute(
"DELETE FROM _catalog_indexes
WHERE schema_name = ?1 AND relation_name = ?2",
params![index.schema, index.name],
)?;
Self::release_relation(&tx, &index, RelationKind::Index)?;
}
drop_fts_aux_tables_for_field(&tx, table_name, column_name)?;
tx.commit()?;
Ok(())
})
}
pub fn rename_column_data(&self, table_name: &str, from: &str, to: &str) -> Result<()> {
let index_updates = self.catalog_index_column_renames(table_name, from, to)?;
self.conn.with_mut(|c| {
let tx = c.savepoint()?;
rename_field_rows_or_keep_existing(
&tx,
"_document_blobs",
"field_name",
table_name,
from,
to,
)?;
for table in [
"_postings",
"_posting_clusters",
"_posting_documents",
"_doc_lengths",
"_field_stats",
"_occurrence_clusters",
"_occurrence_documents",
"_occurrence_lengths",
"_occurrence_fields",
"_vectors",
"_ivf_indexes",
"_ivf_centroids",
"_ivf_assignments",
"_hnsw_indexes",
"_hnsw_nodes",
"_hnsw_edges",
] {
if matches!(
table,
"_postings"
| "_posting_clusters"
| "_posting_documents"
| "_doc_lengths"
| "_field_stats"
) && !Self::table_columns(&tx, table)?
.is_some_and(|columns| columns.contains_key("field"))
{
continue;
}
rename_field_rows_or_keep_existing(&tx, table, "field", table_name, from, to)?;
}
rename_btree_field_rows_or_keep_existing(&tx, table_name, from, to)?;
rename_field_rows_or_keep_existing(
&tx,
"_column_stats",
"column_name",
table_name,
from,
to,
)?;
rename_field_rows_or_keep_existing(
&tx,
"_table_field_analyzers",
"field",
table_name,
from,
to,
)?;
for (index, columns_json) in index_updates {
tx.execute(
"UPDATE _catalog_indexes
SET columns = ?2
WHERE schema_name = ?1 AND relation_name = ?3",
params![index.schema, columns_json, index.name],
)?;
}
rename_fts_aux_tables_for_field(&tx, table_name, from, to)?;
tx.commit()?;
Ok(())
})
}
pub(super) fn catalog_indexes_referencing_column(
&self,
table_name: &str,
column_name: &str,
) -> Result<Vec<RelationIdentity>> {
let mut out = Vec::new();
for row in self.load_catalog_indexes()? {
if row.table_name == table_name
&& columns_json_references(&row.columns_json, column_name)?
{
out.push(row.relation);
}
}
Ok(out)
}
pub(super) fn catalog_index_column_renames(
&self,
table_name: &str,
from: &str,
to: &str,
) -> Result<Vec<(RelationIdentity, String)>> {
let mut out = Vec::new();
for row in self.load_catalog_indexes()? {
if row.table_name != table_name {
continue;
}
if let Some(columns_json) = renamed_columns_json(&row.columns_json, from, to)? {
out.push((row.relation, columns_json));
}
}
Ok(out)
}
}