use super::super::{
params, Catalog, ManagedConnection, OptionalExtension, Result, SQLiteError,
CURRENT_SCHEMA_VERSION,
};
use super::steps::{MigrationAction, MIGRATIONS};
impl Catalog {
pub fn open(conn: ManagedConnection) -> Result<Self> {
let cat = Self::for_initial_restore(conn);
cat.initialize_storage()?;
Ok(cat)
}
pub(crate) fn for_initial_restore(conn: ManagedConnection) -> Self {
Self { conn }
}
pub(in crate::catalog) fn initialize_storage(&self) -> Result<()> {
self.run_migrations()
}
pub fn connection(&self) -> ManagedConnection {
self.conn.clone()
}
pub(super) fn run_migrations(&self) -> Result<()> {
self.conn.with_mut(|conn| {
let mut conn = conn.savepoint()?;
let legacy_meta_only: bool = conn
.query_row(
"SELECT \
(SELECT COUNT(*) FROM sqlite_master \
WHERE type='table' AND name='_meta') > 0 \
AND (SELECT COUNT(*) FROM sqlite_master \
WHERE type='table' AND name='_metadata') = 0",
[],
|r| r.get::<_, i64>(0),
)
.optional()?
.is_some_and(|n| n != 0);
if legacy_meta_only {
conn.execute("ALTER TABLE _meta RENAME TO _metadata", [])?;
}
conn.execute(
"CREATE TABLE IF NOT EXISTS _metadata (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
)",
[],
)?;
let current = conn
.query_row(
"SELECT value FROM _metadata WHERE key = 'schema_version'",
[],
|r| r.get::<_, String>(0),
)
.optional()?;
let current = match current {
Some(version) => version
.parse::<u32>()
.map_err(|_| SQLiteError::InvalidSchemaVersion(version))?,
None => 0,
};
if current > CURRENT_SCHEMA_VERSION {
return Err(SQLiteError::UnsupportedSchemaVersion {
found: current,
supported: CURRENT_SCHEMA_VERSION,
});
}
debug_assert_eq!(
MIGRATIONS.last().map(|migration| migration.version),
Some(CURRENT_SCHEMA_VERSION)
);
for migration in &MIGRATIONS {
if migration.version > current {
let tx = conn.savepoint()?;
match migration.action {
MigrationAction::Sql(sql) => tx.execute_batch(sql)?,
MigrationAction::Custom(migrate) => migrate(&tx)?,
}
tx.execute(
"INSERT OR REPLACE INTO _metadata (key, value) \
VALUES ('schema_version', ?1)",
params![migration.version.to_string()],
)?;
tx.commit()?;
}
}
let repair = conn.savepoint()?;
let schema_before_repair: i64 =
repair.pragma_query_value(None, "schema_version", |row| row.get(0))?;
Self::ensure_column_stats_shape(&repair)?;
let schema_after_repair: i64 =
repair.pragma_query_value(None, "schema_version", |row| row.get(0))?;
if schema_before_repair != schema_after_repair {
Self::install_cache_revision_tracking(&repair)?;
}
repair.commit()?;
conn.commit()?;
Ok(())
})
}
}