use rusqlite::{Transaction, params};
use serde::{Deserialize, Serialize};
const CREATE_HISTORY_TABLE_SQL: &str = "
CREATE TABLE IF NOT EXISTS _row_history (
id INTEGER PRIMARY KEY AUTOINCREMENT,
table_name TEXT NOT NULL,
row_id TEXT NOT NULL,
version INTEGER NOT NULL,
recorded_at INTEGER NOT NULL,
op TEXT NOT NULL CHECK (op IN ('create','update','delete')),
data_json TEXT,
prev_data_json TEXT
)
";
const CREATE_HISTORY_IDX_VERSION_SQL: &str = "
CREATE UNIQUE INDEX IF NOT EXISTS idx_row_history_version
ON _row_history(table_name, row_id, version)
";
const CREATE_HISTORY_IDX_LOOKUP_SQL: &str = "
CREATE INDEX IF NOT EXISTS idx_row_history_lookup
ON _row_history(table_name, row_id, recorded_at DESC)
";
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum HistoryOp {
Create,
Update,
Delete,
}
impl HistoryOp {
pub fn as_str(self) -> &'static str {
match self {
HistoryOp::Create => "create",
HistoryOp::Update => "update",
HistoryOp::Delete => "delete",
}
}
}
impl rusqlite::types::FromSql for HistoryOp {
fn column_result(value: rusqlite::types::ValueRef<'_>) -> rusqlite::types::FromSqlResult<Self> {
let s = String::column_result(value)?;
match s.as_str() {
"create" => Ok(HistoryOp::Create),
"update" => Ok(HistoryOp::Update),
"delete" => Ok(HistoryOp::Delete),
_ => Err(rusqlite::types::FromSqlError::Other(
format!("unknown HistoryOp: {s}").into(),
)),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HistoryRecord {
pub id: i64,
pub table_name: String,
pub row_id: String,
pub version: i64,
pub recorded_at: i64,
pub op: HistoryOp,
pub data_json: Option<String>,
pub prev_data_json: Option<String>,
}
pub fn ensure_history_table(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
conn.execute_batch(CREATE_HISTORY_TABLE_SQL)?;
conn.execute_batch(CREATE_HISTORY_IDX_VERSION_SQL)?;
conn.execute_batch(CREATE_HISTORY_IDX_LOOKUP_SQL)?;
Ok(())
}
pub fn record_in_tx(
tx: &Transaction<'_>,
table_name: &str,
row_id: &str,
op: HistoryOp,
data_json: Option<&str>,
prev_data_json: Option<&str>,
recorded_at: i64,
) -> Result<(), rusqlite::Error> {
let version: i64 = tx
.query_row(
"SELECT COALESCE(MAX(version), 0) + 1 \
FROM _row_history \
WHERE table_name = ?1 AND row_id = ?2",
params![table_name, row_id],
|row| row.get(0),
)
.unwrap_or(1);
tx.execute(
"INSERT INTO _row_history \
(table_name, row_id, version, recorded_at, op, data_json, prev_data_json) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
params![
table_name,
row_id,
version,
recorded_at,
op.as_str(),
data_json,
prev_data_json,
],
)?;
Ok(())
}
pub fn fetch_at(
conn: &rusqlite::Connection,
table_name: &str,
row_id: &str,
at_unix_secs: i64,
) -> Result<Option<HistoryRecord>, rusqlite::Error> {
use rusqlite::OptionalExtension;
conn.query_row(
"SELECT id, table_name, row_id, version, recorded_at, op, data_json, prev_data_json \
FROM _row_history \
WHERE table_name = ?1 AND row_id = ?2 AND recorded_at <= ?3 \
ORDER BY recorded_at DESC, version DESC \
LIMIT 1",
params![table_name, row_id, at_unix_secs],
|row| {
Ok(HistoryRecord {
id: row.get(0)?,
table_name: row.get(1)?,
row_id: row.get(2)?,
version: row.get(3)?,
recorded_at: row.get(4)?,
op: row.get(5)?,
data_json: row.get(6)?,
prev_data_json: row.get(7)?,
})
},
)
.optional()
}
pub fn list_versions(
conn: &rusqlite::Connection,
table_name: &str,
row_id: &str,
) -> Result<Vec<HistoryRecord>, rusqlite::Error> {
let mut stmt = conn.prepare(
"SELECT id, table_name, row_id, version, recorded_at, op, data_json, prev_data_json \
FROM _row_history \
WHERE table_name = ?1 AND row_id = ?2 \
ORDER BY version ASC",
)?;
let rows = stmt.query_map(params![table_name, row_id], |row| {
Ok(HistoryRecord {
id: row.get(0)?,
table_name: row.get(1)?,
row_id: row.get(2)?,
version: row.get(3)?,
recorded_at: row.get(4)?,
op: row.get(5)?,
data_json: row.get(6)?,
prev_data_json: row.get(7)?,
})
})?;
rows.collect()
}
pub fn purge_old_history(
conn: &rusqlite::Connection,
retention_days: u32,
max_per_row: u32,
now_secs: i64,
) -> Result<(), rusqlite::Error> {
if retention_days > 0 {
let cutoff = now_secs - (retention_days as i64) * 86_400;
conn.execute(
"DELETE FROM _row_history WHERE recorded_at < ?1",
params![cutoff],
)?;
}
if max_per_row > 0 {
conn.execute(
"DELETE FROM _row_history \
WHERE id IN ( \
SELECT h.id \
FROM _row_history h \
WHERE ( \
SELECT COUNT(*) \
FROM _row_history h2 \
WHERE h2.table_name = h.table_name \
AND h2.row_id = h.row_id \
AND h2.version >= h.version \
) > ?1 \
)",
params![max_per_row],
)?;
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub fn record_in_tx_for_test(
tx: &rusqlite::Transaction<'_>,
table_name: &str,
row_id: &str,
op: HistoryOp,
data_json: &str,
prev_data_json: Option<&str>,
recorded_at: i64,
version: i64,
) -> Result<(), rusqlite::Error> {
tx.execute(
"INSERT INTO _row_history \
(table_name, row_id, version, recorded_at, op, data_json, prev_data_json) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
rusqlite::params![
table_name,
row_id,
version,
recorded_at,
op.as_str(),
data_json,
prev_data_json,
],
)?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn open_mem() -> rusqlite::Connection {
let conn = rusqlite::Connection::open_in_memory().unwrap();
ensure_history_table(&conn).unwrap();
conn
}
#[test]
fn ensure_table_is_idempotent() {
let conn = open_mem();
ensure_history_table(&conn).unwrap();
}
#[test]
fn record_and_list_versions() {
let conn = open_mem();
let tx = conn.unchecked_transaction().unwrap();
record_in_tx(
&tx,
"t",
"r1",
HistoryOp::Create,
Some("{\"a\":1}"),
None,
1000,
)
.unwrap();
record_in_tx(
&tx,
"t",
"r1",
HistoryOp::Update,
Some("{\"a\":2}"),
Some("{\"a\":1}"),
2000,
)
.unwrap();
tx.commit().unwrap();
let versions = list_versions(&conn, "t", "r1").unwrap();
assert_eq!(versions.len(), 2);
assert_eq!(versions[0].version, 1);
assert_eq!(versions[0].op, HistoryOp::Create);
assert_eq!(versions[1].version, 2);
assert_eq!(versions[1].op, HistoryOp::Update);
}
#[test]
fn fetch_at_returns_correct_snapshot() {
let conn = open_mem();
{
let tx = conn.unchecked_transaction().unwrap();
record_in_tx(
&tx,
"t",
"r1",
HistoryOp::Create,
Some("{\"v\":1}"),
None,
1000,
)
.unwrap();
record_in_tx(
&tx,
"t",
"r1",
HistoryOp::Update,
Some("{\"v\":2}"),
Some("{\"v\":1}"),
2000,
)
.unwrap();
tx.commit().unwrap();
}
let h = fetch_at(&conn, "t", "r1", 1000).unwrap().unwrap();
assert_eq!(h.version, 1);
assert_eq!(h.data_json.as_deref(), Some("{\"v\":1}"));
let h = fetch_at(&conn, "t", "r1", 1500).unwrap().unwrap();
assert_eq!(h.version, 1);
let h = fetch_at(&conn, "t", "r1", 2000).unwrap().unwrap();
assert_eq!(h.version, 2);
let h = fetch_at(&conn, "t", "r1", 999).unwrap();
assert!(h.is_none());
}
#[test]
fn purge_by_age() {
let conn = open_mem();
{
let tx = conn.unchecked_transaction().unwrap();
record_in_tx(&tx, "t", "r1", HistoryOp::Create, Some("{}"), None, 1000).unwrap();
let recent = 1000 + 30 * 86_400 + 1;
record_in_tx(
&tx,
"t",
"r1",
HistoryOp::Update,
Some("{\"a\":1}"),
Some("{}"),
recent,
)
.unwrap();
tx.commit().unwrap();
}
let now = 1000 + 40 * 86_400_i64;
purge_old_history(&conn, 30, 0, now).unwrap();
let versions = list_versions(&conn, "t", "r1").unwrap();
assert_eq!(versions.len(), 1);
assert_eq!(versions[0].op, HistoryOp::Update);
}
#[test]
fn purge_by_max_per_row() {
let conn = open_mem();
{
let tx = conn.unchecked_transaction().unwrap();
for i in 0..5_i64 {
record_in_tx(
&tx,
"t",
"r1",
HistoryOp::Update,
Some("{}"),
Some("{}"),
1000 + i * 100,
)
.unwrap();
}
tx.commit().unwrap();
}
purge_old_history(&conn, 0, 3, 9999).unwrap();
let versions = list_versions(&conn, "t", "r1").unwrap();
assert_eq!(versions.len(), 3, "only newest 3 should remain");
assert_eq!(versions[0].version, 3);
assert_eq!(versions[2].version, 5);
}
}