#[cfg(feature = "change-tracking")]
use std::path::Path;
#[cfg(feature = "change-tracking")]
use oxisql_core::{Row, ToSqlValue};
#[cfg(feature = "change-tracking")]
use oxisql_sqlite_compat::blocking::SqliteConnectionBlocking;
#[cfg(feature = "change-tracking")]
use crate::error::GpkgError;
#[cfg(feature = "change-tracking")]
fn validate_identifier(name: &str) -> Result<(), GpkgError> {
if name.is_empty() {
return Err(GpkgError::ChangeTrackingError(
"identifier must not be empty".to_owned(),
));
}
let mut chars = name.chars();
let first = match chars.next() {
Some(c) => c,
None => {
return Err(GpkgError::ChangeTrackingError(
"identifier must not be empty".to_owned(),
));
}
};
if !first.is_ascii_alphabetic() && first != '_' {
return Err(GpkgError::ChangeTrackingError(format!(
"identifier '{name}' must start with a letter or underscore"
)));
}
for ch in chars {
if !ch.is_ascii_alphanumeric() && ch != '_' {
return Err(GpkgError::ChangeTrackingError(format!(
"identifier '{name}' contains invalid character '{ch}'"
)));
}
}
Ok(())
}
#[cfg(feature = "change-tracking")]
fn parse_change_log_row(row: &Row) -> Result<ChangeLogEntry, GpkgError> {
let id = row
.try_get_by_index::<i64>(0)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
let table_name = row
.try_get_by_index::<String>(1)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
let op_int = row
.try_get_by_index::<i64>(2)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
let feature_id = row
.try_get_by_index::<i64>(3)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
let committed_at = row
.try_get_by_index::<String>(4)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
let operation = ChangeOperation::from_int(op_int)?;
Ok(ChangeLogEntry {
id,
table_name,
operation,
feature_id,
committed_at,
})
}
#[cfg(feature = "change-tracking")]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ChangeOperation {
Insert = 1,
Update = 2,
Delete = 3,
}
#[cfg(feature = "change-tracking")]
impl ChangeOperation {
pub fn from_int(n: i64) -> Result<Self, GpkgError> {
match n {
1 => Ok(Self::Insert),
2 => Ok(Self::Update),
3 => Ok(Self::Delete),
other => Err(GpkgError::ChangeTrackingError(format!(
"unknown operation code {other}"
))),
}
}
pub fn as_int(self) -> i64 {
self as i64
}
}
#[cfg(feature = "change-tracking")]
#[derive(Debug, Clone)]
pub struct ChangeLogEntry {
pub id: i64,
pub table_name: String,
pub operation: ChangeOperation,
pub feature_id: i64,
pub committed_at: String,
}
#[cfg(feature = "change-tracking")]
pub struct ChangeTracker {
conn: SqliteConnectionBlocking,
}
#[cfg(feature = "change-tracking")]
impl ChangeTracker {
pub fn open<P: AsRef<Path>>(path: P) -> Result<Self, GpkgError> {
let path_str = path
.as_ref()
.to_str()
.ok_or_else(|| GpkgError::ChangeTrackingError("path is not valid UTF-8".to_owned()))?;
let conn = SqliteConnectionBlocking::open(path_str)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
Ok(Self { conn })
}
pub fn open_in_memory() -> Result<Self, GpkgError> {
let conn = SqliteConnectionBlocking::open_memory()
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
Ok(Self { conn })
}
pub fn connection(&self) -> &SqliteConnectionBlocking {
&self.conn
}
pub fn create_changes_table(&self) -> Result<(), GpkgError> {
self.conn
.execute_batch(
"CREATE TABLE IF NOT EXISTS gpkg_changes (
id INTEGER PRIMARY KEY AUTOINCREMENT,
table_name TEXT NOT NULL,
operation INTEGER NOT NULL,
feature_id INTEGER NOT NULL,
committed_at TEXT NOT NULL DEFAULT (datetime('now'))
);",
)
.map(|_| ())
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))
}
pub fn enable_tracking(&self, table: &str, fid_column: &str) -> Result<(), GpkgError> {
validate_identifier(table)?;
validate_identifier(fid_column)?;
self.create_changes_table()?;
let insert_ddl = format!(
"CREATE TRIGGER IF NOT EXISTS gpkg_track_{table}_insert \
AFTER INSERT ON {table} \
BEGIN \
INSERT INTO gpkg_changes (table_name, operation, feature_id) \
VALUES ('{table}', 1, NEW.{fid_column}); \
END"
);
let update_ddl = format!(
"CREATE TRIGGER IF NOT EXISTS gpkg_track_{table}_update \
AFTER UPDATE ON {table} \
BEGIN \
INSERT INTO gpkg_changes (table_name, operation, feature_id) \
VALUES ('{table}', 2, NEW.{fid_column}); \
END"
);
let delete_ddl = format!(
"CREATE TRIGGER IF NOT EXISTS gpkg_track_{table}_delete \
AFTER DELETE ON {table} \
BEGIN \
INSERT INTO gpkg_changes (table_name, operation, feature_id) \
VALUES ('{table}', 3, OLD.{fid_column}); \
END"
);
self.conn
.execute(&insert_ddl, &[])
.map(|_| ())
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
self.conn
.execute(&update_ddl, &[])
.map(|_| ())
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
self.conn
.execute(&delete_ddl, &[])
.map(|_| ())
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))
}
pub fn disable_tracking(&self, table: &str) -> Result<(), GpkgError> {
validate_identifier(table)?;
for suffix in ["insert", "update", "delete"] {
let ddl = format!("DROP TRIGGER IF EXISTS gpkg_track_{table}_{suffix}");
self.conn
.execute(&ddl, &[])
.map(|_| ())
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
}
Ok(())
}
pub fn is_tracking(&self, table: &str) -> Result<bool, GpkgError> {
validate_identifier(table)?;
let trigger_name = format!("gpkg_track_{table}_insert");
let rows = self
.conn
.query(
"SELECT COUNT(*) FROM sqlite_master WHERE type='trigger' AND name=$1",
&[&trigger_name as &dyn ToSqlValue],
)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
let count = rows
.first()
.ok_or_else(|| {
GpkgError::ChangeTrackingError("no row returned from COUNT(*) query".to_owned())
})?
.try_get_by_index::<i64>(0)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
Ok(count > 0)
}
pub fn get_all_changes(&self, table: &str) -> Result<Vec<ChangeLogEntry>, GpkgError> {
let rows = self
.conn
.query(
"SELECT id, table_name, operation, feature_id, committed_at
FROM gpkg_changes
WHERE table_name = $1
ORDER BY id ASC",
&[&table as &dyn ToSqlValue],
)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
rows.iter().map(parse_change_log_row).collect()
}
pub fn get_changes_since(
&self,
table: &str,
since_id: i64,
) -> Result<Vec<ChangeLogEntry>, GpkgError> {
let rows = self
.conn
.query(
"SELECT id, table_name, operation, feature_id, committed_at
FROM gpkg_changes
WHERE table_name = $1 AND id > $2
ORDER BY id ASC",
&[&table as &dyn ToSqlValue, &since_id as &dyn ToSqlValue],
)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
rows.iter().map(parse_change_log_row).collect()
}
pub fn clear_changes(&self, table: &str) -> Result<usize, GpkgError> {
let rows_affected = self
.conn
.execute(
"DELETE FROM gpkg_changes WHERE table_name = $1",
&[&table as &dyn ToSqlValue],
)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
Ok(rows_affected as usize)
}
pub fn clear_all_changes(&self) -> Result<usize, GpkgError> {
let rows_affected = self
.conn
.execute("DELETE FROM gpkg_changes", &[])
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
Ok(rows_affected as usize)
}
pub fn tracked_tables(&self) -> Result<Vec<String>, GpkgError> {
let rows = self
.conn
.query(
"SELECT name FROM sqlite_master
WHERE type = 'trigger' AND name LIKE 'gpkg_track_%_insert'",
&[],
)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
rows.iter()
.map(|row| {
let trigger_name = row
.try_get_by_index::<String>(0)
.map_err(|e| GpkgError::ChangeTrackingError(e.to_string()))?;
let inner = trigger_name
.strip_prefix("gpkg_track_")
.and_then(|s| s.strip_suffix("_insert"))
.ok_or_else(|| {
GpkgError::ChangeTrackingError(format!(
"unexpected trigger name format: {trigger_name}"
))
})?;
Ok(inner.to_owned())
})
.collect()
}
}