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)
";
const CREATE_HISTORY_ARCHIVE_TABLE_SQL: &str = "
CREATE TABLE IF NOT EXISTS _row_history_archive (
id INTEGER PRIMARY KEY AUTOINCREMENT,
table_name TEXT NOT NULL,
row_id TEXT NOT NULL,
version_start INTEGER NOT NULL,
version_end INTEGER NOT NULL,
recorded_start INTEGER NOT NULL,
recorded_end INTEGER NOT NULL,
entry_count INTEGER NOT NULL,
format TEXT NOT NULL DEFAULT 'jsonl-zstd',
blob BLOB NOT NULL
)
";
const CREATE_HISTORY_ARCHIVE_IDX_SQL: &str = "
CREATE INDEX IF NOT EXISTS idx_row_history_archive_lookup
ON _row_history_archive(table_name, row_id, recorded_start DESC)
";
const ARCHIVE_ZSTD_LEVEL: i32 = 3;
fn env_tunable(name: &str, default: u32, cell: &'static std::sync::OnceLock<u32>) -> u32 {
*cell.get_or_init(|| {
std::env::var(name)
.ok()
.and_then(|v| v.parse::<u32>().ok())
.unwrap_or(default)
})
}
fn keep_recent() -> u32 {
static CELL: std::sync::OnceLock<u32> = std::sync::OnceLock::new();
env_tunable("MINI_APP_HISTORY_KEEP_RECENT", 16, &CELL)
}
fn chunk_min() -> u32 {
static CELL: std::sync::OnceLock<u32> = std::sync::OnceLock::new();
env_tunable("MINI_APP_HISTORY_CHUNK_MIN", 48, &CELL)
}
#[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)?;
conn.execute_batch(CREATE_HISTORY_ARCHIVE_TABLE_SQL)?;
conn.execute_batch(CREATE_HISTORY_ARCHIVE_IDX_SQL)?;
Ok(())
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct ArchivedEntry {
version: i64,
recorded_at: i64,
op: HistoryOp,
#[serde(skip_serializing_if = "Option::is_none", default)]
data_json: Option<String>,
#[serde(skip_serializing_if = "Option::is_none", default)]
prev_data_json: Option<String>,
}
impl ArchivedEntry {
fn into_record(self, table_name: &str, row_id: &str) -> HistoryRecord {
HistoryRecord {
id: 0,
table_name: table_name.to_string(),
row_id: row_id.to_string(),
version: self.version,
recorded_at: self.recorded_at,
op: self.op,
data_json: self.data_json,
prev_data_json: self.prev_data_json,
}
}
}
fn decode_chunk(blob: &[u8]) -> Result<Vec<ArchivedEntry>, rusqlite::Error> {
let jsonl = zstd::decode_all(blob).map_err(|e| {
rusqlite::Error::ToSqlConversionFailure(
format!("archive chunk zstd decode failed: {e}").into(),
)
})?;
let text = String::from_utf8(jsonl).map_err(|e| {
rusqlite::Error::ToSqlConversionFailure(format!("archive chunk not UTF-8: {e}").into())
})?;
text.lines()
.filter(|l| !l.is_empty())
.map(|l| {
serde_json::from_str::<ArchivedEntry>(l).map_err(|e| {
rusqlite::Error::ToSqlConversionFailure(
format!("archive chunk JSONL parse failed: {e}").into(),
)
})
})
.collect()
}
pub fn record_in_tx(
tx: &Transaction<'_>,
table_name: &str,
row_id: &str,
op: HistoryOp,
data_json: Option<&str>,
recorded_at: i64,
) -> Result<(), rusqlite::Error> {
let version: i64 = tx.query_row(
"SELECT MAX(v) FROM ( \
SELECT COALESCE(MAX(version), 0) AS v \
FROM _row_history \
WHERE table_name = ?1 AND row_id = ?2 \
UNION ALL \
SELECT COALESCE(MAX(version_end), 0) AS v \
FROM _row_history_archive \
WHERE table_name = ?1 AND row_id = ?2 \
)",
params![table_name, row_id],
|row| row.get::<_, i64>(0),
)? + 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, NULL)",
params![
table_name,
row_id,
version,
recorded_at,
op.as_str(),
data_json
],
)?;
maybe_roll_into_archive(
tx,
table_name,
row_id,
keep_recent() as i64,
chunk_min() as i64,
)?;
Ok(())
}
fn maybe_roll_into_archive(
tx: &Transaction<'_>,
table_name: &str,
row_id: &str,
keep: i64,
chunk_min: i64,
) -> Result<(), rusqlite::Error> {
let count: i64 = tx.query_row(
"SELECT COUNT(*) FROM _row_history WHERE table_name = ?1 AND row_id = ?2",
params![table_name, row_id],
|row| row.get(0),
)?;
if count <= keep + chunk_min {
return Ok(());
}
let roll_n = (count - keep).min(chunk_min + 1);
let mut stmt = tx.prepare(
"SELECT 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 \
LIMIT ?3",
)?;
let mut ids: Vec<i64> = Vec::with_capacity(roll_n as usize);
let mut entries: Vec<ArchivedEntry> = Vec::with_capacity(roll_n as usize);
let rows = stmt.query_map(params![table_name, row_id, roll_n], |row| {
Ok((
row.get::<_, i64>(0)?,
ArchivedEntry {
version: row.get(1)?,
recorded_at: row.get(2)?,
op: row.get(3)?,
data_json: row.get(4)?,
prev_data_json: row.get(5)?,
},
))
})?;
for r in rows {
let (id, entry) = r?;
ids.push(id);
entries.push(entry);
}
drop(stmt);
if entries.is_empty() {
return Ok(());
}
let mut jsonl = String::new();
for e in &entries {
jsonl.push_str(
&serde_json::to_string(e).expect("ArchivedEntry serialization is infallible"),
);
jsonl.push('\n');
}
let blob = zstd::encode_all(jsonl.as_bytes(), ARCHIVE_ZSTD_LEVEL).map_err(|e| {
rusqlite::Error::ToSqlConversionFailure(
format!("archive chunk zstd encode failed: {e}").into(),
)
})?;
let first = entries.first().expect("non-empty");
let last = entries.last().expect("non-empty");
tx.execute(
"INSERT INTO _row_history_archive \
(table_name, row_id, version_start, version_end, \
recorded_start, recorded_end, entry_count, format, blob) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, 'jsonl-zstd', ?8)",
params![
table_name,
row_id,
first.version,
last.version,
first.recorded_at,
last.recorded_at,
entries.len() as i64,
blob,
],
)?;
for chunk in ids.chunks(500) {
let placeholders = vec!["?"; chunk.len()].join(",");
let sql = format!("DELETE FROM _row_history WHERE id IN ({placeholders})");
tx.execute(&sql, rusqlite::params_from_iter(chunk.iter()))?;
}
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;
let raw = 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()?;
if raw.is_some() {
return Ok(raw);
}
let blob: Option<Vec<u8>> = conn
.query_row(
"SELECT blob FROM _row_history_archive \
WHERE table_name = ?1 AND row_id = ?2 AND recorded_start <= ?3 \
ORDER BY recorded_start DESC, version_end DESC \
LIMIT 1",
params![table_name, row_id, at_unix_secs],
|row| row.get(0),
)
.optional()?;
let Some(blob) = blob else {
return Ok(None);
};
let entries = decode_chunk(&blob)?;
Ok(entries
.into_iter()
.filter(|e| e.recorded_at <= at_unix_secs)
.max_by_key(|e| (e.recorded_at, e.version))
.map(|e| e.into_record(table_name, row_id)))
}
pub fn list_versions(
conn: &rusqlite::Connection,
table_name: &str,
row_id: &str,
) -> Result<Vec<HistoryRecord>, rusqlite::Error> {
let mut out: Vec<HistoryRecord> = Vec::new();
let mut astmt = conn.prepare(
"SELECT blob FROM _row_history_archive \
WHERE table_name = ?1 AND row_id = ?2 \
ORDER BY version_start ASC",
)?;
let blobs = astmt.query_map(params![table_name, row_id], |row| row.get::<_, Vec<u8>>(0))?;
for blob in blobs {
for e in decode_chunk(&blob?)? {
out.push(e.into_record(table_name, row_id));
}
}
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)?,
})
})?;
for r in rows {
out.push(r?);
}
Ok(out)
}
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}"), 1000).unwrap();
record_in_tx(&tx, "t", "r1", HistoryOp::Update, Some("{\"a\":2}"), 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}"), 1000).unwrap();
record_in_tx(&tx, "t", "r1", HistoryOp::Update, Some("{\"v\":2}"), 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("{}"), 1000).unwrap();
let recent = 1000 + 30 * 86_400 + 1;
record_in_tx(&tx, "t", "r1", HistoryOp::Update, Some("{\"a\":1}"), 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 roll_moves_oldest_entries_into_archive_and_reads_through() {
let conn = open_mem();
{
let tx = conn.unchecked_transaction().unwrap();
for i in 0..10_i64 {
let body = format!("{{\"v\":{}}}", i + 1);
record_in_tx(
&tx,
"t",
"r1",
HistoryOp::Update,
Some(&body),
1000 + i * 100,
)
.unwrap();
}
maybe_roll_into_archive(&tx, "t", "r1", 3, 2).unwrap();
maybe_roll_into_archive(&tx, "t", "r1", 3, 2).unwrap();
maybe_roll_into_archive(&tx, "t", "r1", 3, 2).unwrap();
tx.commit().unwrap();
}
let raw_count: i64 = conn
.query_row("SELECT COUNT(*) FROM _row_history", [], |r| r.get(0))
.unwrap();
assert_eq!(raw_count, 4);
let chunks: Vec<(i64, i64, i64)> = conn
.prepare(
"SELECT version_start, version_end, entry_count \
FROM _row_history_archive ORDER BY version_start",
)
.unwrap()
.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))
.unwrap()
.collect::<Result<_, _>>()
.unwrap();
assert_eq!(chunks, vec![(1, 3, 3), (4, 6, 3)]);
let versions = list_versions(&conn, "t", "r1").unwrap();
assert_eq!(versions.len(), 10);
assert_eq!(
versions.iter().map(|v| v.version).collect::<Vec<_>>(),
(1..=10).collect::<Vec<i64>>()
);
assert!(versions[..6].iter().all(|v| v.id == 0));
assert!(versions[6..].iter().all(|v| v.id > 0));
let h = fetch_at(&conn, "t", "r1", 1450).unwrap().unwrap();
assert_eq!(h.version, 5);
assert_eq!(h.data_json.as_deref(), Some("{\"v\":5}"));
let h = fetch_at(&conn, "t", "r1", 1900).unwrap().unwrap();
assert_eq!(h.version, 10);
assert!(h.id > 0);
assert!(fetch_at(&conn, "t", "r1", 999).unwrap().is_none());
}
#[test]
fn roll_is_capped_per_call_for_backlog_drain() {
let conn = open_mem();
let tx = conn.unchecked_transaction().unwrap();
for i in 0..100_i64 {
record_in_tx_for_test(
&tx,
"t",
"r1",
HistoryOp::Update,
"{}",
None,
1000 + i,
i + 1,
)
.unwrap();
}
maybe_roll_into_archive(&tx, "t", "r1", 3, 2).unwrap();
tx.commit().unwrap();
let archived: i64 = conn
.query_row(
"SELECT COALESCE(SUM(entry_count), 0) FROM _row_history_archive",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(archived, 3);
let raw: i64 = conn
.query_row("SELECT COUNT(*) FROM _row_history", [], |r| r.get(0))
.unwrap();
assert_eq!(raw, 97);
}
#[test]
fn version_stays_monotonic_across_roll() {
let conn = open_mem();
let tx = conn.unchecked_transaction().unwrap();
for i in 0..10_i64 {
record_in_tx(&tx, "t", "r1", HistoryOp::Update, Some("{}"), 1000 + i).unwrap();
}
maybe_roll_into_archive(&tx, "t", "r1", 3, 2).unwrap();
record_in_tx(&tx, "t", "r1", HistoryOp::Update, Some("{}"), 2000).unwrap();
tx.commit().unwrap();
let versions = list_versions(&conn, "t", "r1").unwrap();
assert_eq!(versions.last().unwrap().version, 11);
}
#[test]
fn roll_below_threshold_is_noop() {
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("{}"), 1000 + i).unwrap();
}
maybe_roll_into_archive(&tx, "t", "r1", 3, 2).unwrap();
tx.commit().unwrap();
let archive_count: i64 = conn
.query_row("SELECT COUNT(*) FROM _row_history_archive", [], |r| {
r.get(0)
})
.unwrap();
assert_eq!(archive_count, 0, "5 <= keep(3)+chunk_min(2) must not roll");
}
#[test]
fn prev_data_json_is_not_written() {
let conn = open_mem();
let tx = conn.unchecked_transaction().unwrap();
record_in_tx(&tx, "t", "r1", HistoryOp::Create, Some("{\"a\":1}"), 1000).unwrap();
record_in_tx(&tx, "t", "r1", HistoryOp::Update, Some("{\"a\":2}"), 2000).unwrap();
tx.commit().unwrap();
let n: i64 = conn
.query_row(
"SELECT COUNT(*) FROM _row_history WHERE prev_data_json IS NOT NULL",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(n, 0);
}
#[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("{}"),
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);
}
}