1use 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 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 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 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}