use super::{
migration_relation, params, Catalog, OptionalExtension, RelationIdentity, RelationKind, Result,
SQLiteError, SequenceOptions, SequenceReservationResult, SequenceRow, ViewRow,
};
use uqa_storage::catalog::{
sequence_value_reservation, SequenceOwner, SequenceOwnerDependency, SequenceValuePosition,
};
fn concrete_sequence_options(sequence: &SequenceRow) -> SequenceOptions {
let default_min = if sequence.increment > 0 { 1 } else { i64::MIN };
let default_max = if sequence.increment > 0 { i64::MAX } else { -1 };
SequenceOptions {
data_type: sequence.options.data_type.clone(),
min_value: Some(sequence.options.min_value.unwrap_or(default_min)),
max_value: Some(sequence.options.max_value.unwrap_or(default_max)),
cycle: sequence.options.cycle,
cache_size: sequence.options.cache_size,
}
}
fn decode_sequence_owner(
relation: &RelationIdentity,
table_object_id: Option<Vec<u8>>,
column_object_id: Option<Vec<u8>>,
dependency: Option<String>,
) -> Result<Option<SequenceOwner>> {
let (table_object_id, column_object_id, dependency) =
match (table_object_id, column_object_id, dependency) {
(None, None, None) => return Ok(None),
(Some(table), Some(column), Some(dependency)) => (table, column, dependency),
_ => {
return Err(SQLiteError::StorageBackend(format!(
"corrupt sequence `{}` has an incomplete owner dependency",
relation.qualified_name()
)))
}
};
let table_object_id: [u8; 16] = table_object_id.try_into().map_err(|value: Vec<u8>| {
SQLiteError::StorageBackend(format!(
"corrupt sequence `{}` owner table identity has {} bytes",
relation.qualified_name(),
value.len()
))
})?;
let column_object_id: [u8; 16] = column_object_id.try_into().map_err(|value: Vec<u8>| {
SQLiteError::StorageBackend(format!(
"corrupt sequence `{}` owner column identity has {} bytes",
relation.qualified_name(),
value.len()
))
})?;
let dependency = match dependency.as_str() {
"a" => SequenceOwnerDependency::Automatic,
"i" => SequenceOwnerDependency::Internal,
other => {
return Err(SQLiteError::StorageBackend(format!(
"corrupt sequence `{}` owner dependency `{other}`",
relation.qualified_name()
)))
}
};
Ok(Some(SequenceOwner {
table_object_id,
column_object_id,
dependency,
}))
}
struct RawSequenceRow {
schema: String,
name: String,
object_id: Vec<u8>,
definition_generation: Vec<u8>,
start: i64,
increment: i64,
current: i64,
called: bool,
persistence: String,
data_type: String,
min_value: i64,
max_value: i64,
cycle: bool,
cache_size: i64,
owner_table_object_id: Option<Vec<u8>>,
owner_column_object_id: Option<Vec<u8>>,
owner_dependency: Option<String>,
role_owner: String,
acl_json: Option<String>,
log_count: i64,
}
fn read_raw_sequence_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<RawSequenceRow> {
Ok(RawSequenceRow {
schema: row.get(0)?,
name: row.get(1)?,
object_id: row.get(2)?,
definition_generation: row.get(3)?,
start: row.get(4)?,
increment: row.get(5)?,
current: row.get(6)?,
called: row.get(7)?,
persistence: row.get(8)?,
data_type: row.get(9)?,
min_value: row.get(10)?,
max_value: row.get(11)?,
cycle: row.get(12)?,
cache_size: row.get(13)?,
owner_table_object_id: row.get(14)?,
owner_column_object_id: row.get(15)?,
owner_dependency: row.get(16)?,
role_owner: row.get(17)?,
acl_json: row.get(18)?,
log_count: row.get(19)?,
})
}
fn decode_sequence_identity(
relation: &RelationIdentity,
label: &str,
value: Vec<u8>,
) -> Result<[u8; 16]> {
value.try_into().map_err(|value: Vec<u8>| {
SQLiteError::StorageBackend(format!(
"corrupt sequence `{}` {label} has {} bytes",
relation.qualified_name(),
value.len()
))
})
}
fn decode_raw_sequence_row(raw: RawSequenceRow) -> Result<SequenceRow> {
let relation = RelationIdentity::new(raw.schema, raw.name);
Ok(SequenceRow {
role_owner: raw.role_owner,
acl: raw
.acl_json
.map(|json| serde_json::from_str(&json))
.transpose()?,
owner: decode_sequence_owner(
&relation,
raw.owner_table_object_id,
raw.owner_column_object_id,
raw.owner_dependency,
)?,
object_id: decode_sequence_identity(&relation, "object identity", raw.object_id)?,
definition_generation: decode_sequence_identity(
&relation,
"definition generation",
raw.definition_generation,
)?,
relation,
start: raw.start,
increment: raw.increment,
current: raw.current,
called: raw.called,
log_count: raw.log_count,
persistence: raw.persistence,
options: SequenceOptions {
data_type: raw.data_type,
min_value: Some(raw.min_value),
max_value: Some(raw.max_value),
cycle: raw.cycle,
cache_size: raw.cache_size,
},
})
}
fn reserve_sequence_values_in_connection(
connection: &rusqlite::Connection,
relation: &RelationIdentity,
object_id: [u8; 16],
definition_generation: [u8; 16],
) -> Result<SequenceReservationResult> {
let stored = connection
.query_row(
"SELECT object_id, definition_generation, current, called, increment, min_value, max_value, cycle, cache_size, log_count
FROM _sequences WHERE schema_name = ?1 AND relation_name = ?2",
params![relation.schema, relation.name],
|row| {
Ok((
row.get::<_, Vec<u8>>(0)?,
row.get::<_, Vec<u8>>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, bool>(3)?,
row.get::<_, i64>(4)?,
row.get::<_, i64>(5)?,
row.get::<_, i64>(6)?,
row.get::<_, bool>(7)?,
row.get::<_, i64>(8)?,
row.get::<_, i64>(9)?,
))
},
)
.optional()?;
let Some((
stored_object_id,
stored_generation,
current,
called,
increment,
min,
max,
cycle,
cache_size,
log_count,
)) = stored
else {
return Ok(SequenceReservationResult::Missing);
};
let stored_object_id: [u8; 16] = stored_object_id.try_into().map_err(|value: Vec<u8>| {
SQLiteError::StorageBackend(format!(
"corrupt sequence `{}` object identity has {} bytes",
relation.qualified_name(),
value.len()
))
})?;
if stored_object_id != object_id {
return Ok(SequenceReservationResult::Missing);
}
let stored_generation: [u8; 16] = stored_generation.try_into().map_err(|value: Vec<u8>| {
SQLiteError::StorageBackend(format!(
"corrupt sequence `{}` definition generation has {} bytes",
relation.qualified_name(),
value.len()
))
})?;
if stored_generation != definition_generation {
return Ok(SequenceReservationResult::DefinitionChanged);
}
if increment == 0 || cache_size <= 0 {
return Err(SQLiteError::StorageBackend(format!(
"corrupt sequence `{}` has increment {increment} and cache size {cache_size}",
relation.qualified_name()
)));
}
let Some(reservation) = sequence_value_reservation(
SequenceValuePosition {
current,
called,
log_count,
},
increment,
min,
max,
cycle,
cache_size,
) else {
return Ok(SequenceReservationResult::Exhausted);
};
let updated = connection.execute(
"UPDATE _sequences SET current = ?5, called = 1, log_count = ?6
WHERE schema_name = ?1 AND relation_name = ?2 AND object_id = ?3 AND definition_generation = ?4",
params![
relation.schema,
relation.name,
object_id.as_slice(),
definition_generation.as_slice(),
reservation.last_value,
reservation.log_count,
],
)?;
if updated != 1 {
return Err(SQLiteError::StorageBackend(format!(
"sequence `{}` changed while reserving cached values",
relation.qualified_name()
)));
}
Ok(SequenceReservationResult::Reserved(reservation))
}
impl Catalog {
pub fn create_sequence_row(&self, sequence: &SequenceRow) -> Result<bool> {
self.conn.with_mut(|connection| {
let tx = connection.savepoint()?;
let exists = tx
.query_row(
"SELECT 1 FROM _sequences
WHERE schema_name = ?1 AND relation_name = ?2",
params![sequence.relation.schema, sequence.relation.name],
|_| Ok(()),
)
.optional()?
.is_some();
if exists {
return Ok(false);
}
Self::claim_relation(&tx, &sequence.relation, RelationKind::Sequence)?;
let options = concrete_sequence_options(sequence);
let owner_table = sequence.owner.map(|owner| owner.table_object_id);
let owner_column = sequence.owner.map(|owner| owner.column_object_id);
let owner_dependency = sequence
.owner
.map(|owner| owner.dependency.catalog_code());
let acl_json = sequence
.acl
.as_ref()
.map(serde_json::to_string)
.transpose()?;
tx.execute(
"INSERT INTO _sequences
(schema_name, relation_name, kind, object_id, definition_generation, start, increment, current, called, persistence, data_type, min_value, max_value, cycle, cache_size, owner_table_object_id, owner_column_object_id, owner_dependency, role_owner, acl_json, log_count)
VALUES (?1, ?2, 'sequence', ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19, ?20)",
params![
sequence.relation.schema,
sequence.relation.name,
sequence.object_id.as_slice(),
sequence.definition_generation.as_slice(),
sequence.start,
sequence.increment,
sequence.current,
sequence.called,
sequence.persistence,
options.data_type,
options.min_value,
options.max_value,
options.cycle,
options.cache_size,
owner_table.as_ref().map(<[u8; 16]>::as_slice),
owner_column.as_ref().map(<[u8; 16]>::as_slice),
owner_dependency,
sequence.role_owner,
acl_json,
sequence.log_count,
],
)?;
tx.commit()?;
Ok(true)
})
}
pub fn replace_sequence_row(&self, sequence: &SequenceRow) -> Result<bool> {
self.conn.with(|connection| {
let options = concrete_sequence_options(sequence);
let owner_table = sequence.owner.map(|owner| owner.table_object_id);
let owner_column = sequence.owner.map(|owner| owner.column_object_id);
let owner_dependency = sequence
.owner
.map(|owner| owner.dependency.catalog_code());
let acl_json = sequence
.acl
.as_ref()
.map(serde_json::to_string)
.transpose()?;
Ok(connection.execute(
"UPDATE _sequences
SET object_id = ?3, definition_generation = ?4, start = ?5, increment = ?6, current = ?7, called = ?8, persistence = ?9,
data_type = ?10, min_value = ?11, max_value = ?12, cycle = ?13, cache_size = ?14,
owner_table_object_id = ?15, owner_column_object_id = ?16, owner_dependency = ?17, role_owner = ?18, acl_json = ?19, log_count = ?20
WHERE schema_name = ?1 AND relation_name = ?2",
params![
sequence.relation.schema,
sequence.relation.name,
sequence.object_id.as_slice(),
sequence.definition_generation.as_slice(),
sequence.start,
sequence.increment,
sequence.current,
sequence.called,
sequence.persistence,
options.data_type,
options.min_value,
options.max_value,
options.cycle,
options.cache_size,
owner_table.as_ref().map(<[u8; 16]>::as_slice),
owner_column.as_ref().map(<[u8; 16]>::as_slice),
owner_dependency,
sequence.role_owner,
acl_json,
sequence.log_count,
],
)? != 0)
})
}
pub fn rename_sequence_row(&self, from: &str, to: &str) -> Result<bool> {
let from_relation = migration_relation(from)?;
let to_relation = migration_relation(to)?;
self.conn.with_mut(|connection| {
let tx = connection.savepoint()?;
let source_exists = tx
.query_row(
"SELECT 1 FROM _sequences
WHERE schema_name = ?1 AND relation_name = ?2",
params![from_relation.schema, from_relation.name],
|_| Ok(()),
)
.optional()?
.is_some();
if !source_exists {
return Ok(false);
}
if from_relation == to_relation {
return Ok(true);
}
let target_kind = tx
.query_row(
"SELECT kind FROM _relations
WHERE schema_name = ?1 AND relation_name = ?2",
params![to_relation.schema, to_relation.name],
|row| row.get::<_, String>(0),
)
.optional()?;
if let Some(kind) = target_kind {
return Err(SQLiteError::StorageBackend(format!(
"relation `{}` already exists as {kind}",
to_relation.qualified_name()
)));
}
Self::claim_relation(&tx, &to_relation, RelationKind::Sequence)?;
let updated = tx.execute(
"UPDATE _sequences
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 != 1 {
return Err(SQLiteError::StorageBackend(format!(
"sequence `{from}` changed while renaming"
)));
}
Self::release_relation(&tx, &from_relation, RelationKind::Sequence)?;
tx.commit()?;
Ok(true)
})
}
pub fn drop_sequence_row(&self, name: &str) -> Result<bool> {
let relation = migration_relation(name)?;
self.conn.with_mut(|connection| {
let tx = connection.savepoint()?;
let removed = tx.execute(
"DELETE FROM _sequences
WHERE schema_name = ?1 AND relation_name = ?2",
params![relation.schema, relation.name],
)? != 0;
if removed {
Self::release_relation(&tx, &relation, RelationKind::Sequence)?;
}
tx.commit()?;
Ok(removed)
})
}
pub fn load_sequence_rows(&self) -> Result<Vec<SequenceRow>> {
self.conn.with(|connection| {
let mut statement = connection.prepare(
"SELECT schema_name, relation_name, object_id, definition_generation, start, increment, current, called, persistence,
data_type, min_value, max_value, cycle, cache_size, owner_table_object_id, owner_column_object_id, owner_dependency, role_owner, acl_json, log_count
FROM _sequences ORDER BY schema_name, relation_name",
)?;
let sequences = statement
.query_map([], read_raw_sequence_row)?
.map(|row| decode_raw_sequence_row(row?))
.collect();
sequences
})
}
pub fn reserve_sequence_values(
&self,
name: &str,
object_id: [u8; 16],
definition_generation: [u8; 16],
) -> Result<SequenceReservationResult> {
let relation = migration_relation(name)?;
self.conn.with_mut(|connection| {
if connection.is_autocommit() {
let tx = connection
.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
let result = reserve_sequence_values_in_connection(
&tx,
&relation,
object_id,
definition_generation,
)?;
tx.commit()?;
return Ok(result);
}
let tx = connection.savepoint()?;
let result = reserve_sequence_values_in_connection(
&tx,
&relation,
object_id,
definition_generation,
)?;
tx.commit()?;
Ok(result)
})
}
pub fn set_sequence_value(
&self,
name: &str,
object_id: [u8; 16],
value: i64,
called: bool,
log_count: i64,
) -> Result<Option<i64>> {
let relation = migration_relation(name)?;
self.conn.with(|connection| {
Ok(connection
.query_row(
"UPDATE _sequences SET current = ?4, called = ?5, log_count = ?6
WHERE schema_name = ?1 AND relation_name = ?2 AND object_id = ?3 RETURNING current",
params![
relation.schema,
relation.name,
object_id.as_slice(),
value,
called,
log_count,
],
|row| row.get(0),
)
.optional()?)
})
}
pub fn save_view(&self, view: &ViewRow) -> Result<()> {
self.conn.with_mut(|connection| {
let tx = connection.savepoint()?;
Self::claim_relation(&tx, &view.relation, RelationKind::View)?;
let acl_json = view.acl.as_ref().map(serde_json::to_string).transpose()?;
let column_acls_json = serde_json::to_string(&view.column_acls)?;
tx.execute(
"INSERT OR REPLACE INTO _views
(schema_name, relation_name, kind, role_owner, acl_json, column_acls_json, definition_json)
VALUES (?1, ?2, 'view', ?3, ?4, ?5, ?6)",
params![
view.relation.schema,
view.relation.name,
view.role_owner,
acl_json,
column_acls_json,
view.definition_json
],
)?;
tx.commit()?;
Ok(())
})
}
pub fn rename_view(&self, from: &RelationIdentity, to: &RelationIdentity) -> Result<bool> {
if from.schema != to.schema {
return Err(SQLiteError::StorageBackend(
"moving a view between schemas is not supported by the catalog".into(),
));
}
self.conn.with_mut(|connection| {
let source_exists = connection.query_row(
"SELECT EXISTS(SELECT 1 FROM _views WHERE schema_name = ?1 AND relation_name = ?2)",
params![from.schema, from.name],
|row| row.get::<_, bool>(0),
)?;
if from == to || !source_exists {
return Ok(source_exists);
}
let target_exists = connection.query_row(
"SELECT EXISTS(SELECT 1 FROM _relations WHERE schema_name = ?1 AND relation_name = ?2)",
params![to.schema, to.name],
|row| row.get::<_, bool>(0),
)?;
if target_exists {
return Err(SQLiteError::StorageBackend(format!(
"relation `{}` already exists",
to.qualified_name()
)));
}
let tx = connection.savepoint()?;
Self::claim_relation(&tx, to, RelationKind::View)?;
let updated = tx.execute(
"UPDATE _views SET schema_name = ?3, relation_name = ?4 WHERE schema_name = ?1 AND relation_name = ?2",
params![from.schema, from.name, to.schema, to.name],
)?;
if updated != 1 {
return Err(SQLiteError::StorageBackend(format!(
"view `{}` disappeared during rename",
from.qualified_name()
)));
}
Self::release_relation(&tx, from, RelationKind::View)?;
tx.commit()?;
Ok(true)
})
}
pub fn drop_view(&self, relation: &RelationIdentity) -> Result<bool> {
self.conn.with_mut(|connection| {
let tx = connection.savepoint()?;
let removed = tx.execute(
"DELETE FROM _views WHERE schema_name = ?1 AND relation_name = ?2",
params![relation.schema, relation.name],
)? != 0;
if removed {
Self::release_relation(&tx, relation, RelationKind::View)?;
}
tx.commit()?;
Ok(removed)
})
}
pub fn load_views(&self) -> Result<Vec<ViewRow>> {
self.conn.with(|connection| {
let mut statement = connection.prepare(
"SELECT schema_name, relation_name, role_owner, acl_json, column_acls_json, definition_json
FROM _views ORDER BY schema_name, relation_name",
)?;
let rows = statement.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, Option<String>>(3)?,
row.get::<_, Option<String>>(4)?,
row.get::<_, String>(5)?,
))
})?;
let mut views = Vec::new();
for row in rows {
let (schema, name, role_owner, acl_json, column_acls_json, definition_json) = row?;
views.push(ViewRow {
relation: RelationIdentity::new(schema, name),
role_owner,
acl: acl_json
.as_deref()
.map(serde_json::from_str)
.transpose()?,
column_acls: column_acls_json
.as_deref()
.map(serde_json::from_str)
.transpose()?
.unwrap_or_default(),
definition_json,
});
}
Ok(views)
})
}
}