use std::path::Path;
use libsql::TransactionBehavior;
use crate::error::{DbError, Result, WriteOp};
use crate::schema::ddl::ARCHIVE_SESSION_MARKER;
use crate::util::limits::HYDRATE_CHUNK;
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct ArchiveReport {
pub links_archived: usize,
pub concepts_archived: usize,
pub log_entries_archived: usize,
pub horizon: Option<i64>,
}
const COLD_SCHEMA: &[&str] = &[
r#"CREATE TABLE IF NOT EXISTS cold.links (
source_id TEXT NOT NULL,
target_id TEXT NOT NULL,
edge_type TEXT NOT NULL,
valid_from TEXT NOT NULL,
recorded_at TEXT NOT NULL,
valid_to TEXT NOT NULL,
weight REAL NOT NULL CHECK (weight >= 0.0 AND weight < 9e999 AND typeof(weight) = 'real'),
properties TEXT NOT NULL,
branch_id TEXT NOT NULL DEFAULT 'main',
PRIMARY KEY (source_id, target_id, edge_type, valid_from, recorded_at, branch_id)
)"#,
r#"CREATE TABLE IF NOT EXISTS cold.concepts (
rowid_pk INTEGER,
id TEXT NOT NULL PRIMARY KEY,
title TEXT NOT NULL,
content TEXT NOT NULL DEFAULT '',
embedding_model TEXT,
valid_from TEXT NOT NULL,
valid_to TEXT NOT NULL,
recorded_at TEXT NOT NULL,
retired INTEGER NOT NULL DEFAULT 0,
branch_id TEXT NOT NULL DEFAULT 'main'
)"#,
r#"CREATE TABLE IF NOT EXISTS cold.transaction_log (
seq_id INTEGER PRIMARY KEY,
table_name TEXT NOT NULL,
entity_id TEXT NOT NULL,
operation TEXT NOT NULL,
payload TEXT NOT NULL,
recorded_at TEXT NOT NULL,
branch_id TEXT NOT NULL DEFAULT 'main'
)"#,
"CREATE INDEX IF NOT EXISTS cold.idx_cold_txlog_entity ON transaction_log (entity_id)",
"CREATE INDEX IF NOT EXISTS cold.idx_cold_txlog_time ON transaction_log (recorded_at)",
r#"CREATE TABLE IF NOT EXISTS cold.branches (
branch_id TEXT NOT NULL PRIMARY KEY,
parent_id TEXT,
forked_at TEXT,
created_at TEXT NOT NULL,
archived_at TEXT NOT NULL
)"#,
r#"CREATE TABLE IF NOT EXISTS cold.archive_horizon (
archived_at TEXT NOT NULL,
cutoff TEXT NOT NULL,
horizon INTEGER
)"#,
];
const LINKS_ARCHIVABLE: &str = r#"
recorded_at < :cutoff AND (
EXISTS (
SELECT 1 FROM links newer
WHERE newer.source_id = links.source_id
AND newer.target_id = links.target_id
AND newer.edge_type = links.edge_type
AND newer.valid_from = links.valid_from
AND newer.branch_id = links.branch_id
AND newer.recorded_at > links.recorded_at
AND NOT EXISTS (
SELECT 1 FROM branches b
WHERE b.forked_at IS NOT NULL
AND b.forked_at >= links.recorded_at
AND b.forked_at < newer.recorded_at
)
)
OR (valid_to <> '9999-12-31T23:59:59.999999Z' AND valid_to <= :cutoff
AND NOT EXISTS (
SELECT 1 FROM links other
WHERE other.source_id = links.source_id
AND other.target_id = links.target_id
AND other.edge_type = links.edge_type
AND other.valid_from = links.valid_from
AND other.branch_id <> links.branch_id
)
AND NOT EXISTS (
SELECT 1 FROM links older, branches b
WHERE older.source_id = links.source_id
AND older.target_id = links.target_id
AND older.edge_type = links.edge_type
AND older.valid_from = links.valid_from
AND older.branch_id = links.branch_id
AND older.recorded_at < links.recorded_at
AND b.forked_at IS NOT NULL
AND b.forked_at >= older.recorded_at
AND b.forked_at < links.recorded_at
))
)
"#;
const LOG_ARCHIVABLE: &str = r#"
recorded_at < :cutoff AND EXISTS (
SELECT 1 FROM transaction_log newer
WHERE newer.entity_id = transaction_log.entity_id
AND newer.branch_id = transaction_log.branch_id
AND newer.seq_id > transaction_log.seq_id
AND NOT EXISTS (
SELECT 1 FROM branches b
WHERE b.forked_at IS NOT NULL
AND b.forked_at >= transaction_log.recorded_at
AND b.forked_at < newer.recorded_at
)
)
"#;
const CONCEPTS_ARCHIVABLE: &str = r#"
retired = 1
AND recorded_at < :cutoff
AND valid_to < :cutoff
AND NOT EXISTS (
SELECT 1 FROM links
WHERE links.source_id = concepts.id
OR links.target_id = concepts.id
)
"#;
pub async fn archivable_concepts(conn: &libsql::Connection, cutoff: &str) -> Result<Vec<String>> {
let mut rows = conn
.query(
&format!("SELECT id FROM concepts WHERE {CONCEPTS_ARCHIVABLE} ORDER BY id"),
libsql::named_params! {":cutoff": cutoff},
)
.await?;
let mut ids = Vec::new();
while let Some(row) = rows.next().await? {
ids.push(row.get::<String>(0)?);
}
Ok(ids)
}
pub(crate) fn archive_present(path: &Path) -> bool {
std::fs::metadata(path).is_ok_and(|m| m.len() > 0)
}
pub async fn archive(
conn: &libsql::Connection,
cutoff: &str,
archived_at: &str,
archive_path: &Path,
) -> Result<ArchiveReport> {
crate::temporal::replay::detach_stale_cold(conn).await;
conn.execute(
"ATTACH DATABASE ?1 AS cold",
libsql::params![archive_path.to_string_lossy().as_ref()],
)
.await?;
let result = archive_session(conn, cutoff, archived_at).await;
if let Err(e) = conn.execute("DETACH DATABASE cold", ()).await {
tracing::warn!("archive: failed to DETACH cold database: {e}");
}
result
}
const COLD_LINEAGE_INDICES: &[&str] = &[
"CREATE INDEX IF NOT EXISTS cold.idx_cold_txlog_fold_partition ON transaction_log \
(table_name, entity_id, branch_id, seq_id DESC)",
];
async fn upgrade_cold_lineage(tx: &libsql::Transaction) -> Result<()> {
for table in ["links", "concepts", "transaction_log"] {
if !cold_has_branch(tx, table).await? {
tx.execute(
&format!(
"ALTER TABLE cold.{table} ADD COLUMN branch_id TEXT NOT NULL DEFAULT 'main'"
),
(),
)
.await?;
}
}
if !cold_links_keyed_by_lineage(tx).await? {
tx.execute(COLD_LINKS_V15, ()).await?;
tx.execute(
"INSERT INTO cold.links_v15 \
(source_id, target_id, edge_type, valid_from, recorded_at, \
valid_to, weight, properties, branch_id) \
SELECT source_id, target_id, edge_type, valid_from, recorded_at, \
valid_to, weight, properties, branch_id FROM cold.links",
(),
)
.await?;
tx.execute("DROP TABLE cold.links", ()).await?;
tx.execute("ALTER TABLE cold.links_v15 RENAME TO links", ())
.await?;
}
for ddl in COLD_LINEAGE_INDICES {
tx.execute(ddl, ()).await?;
}
Ok(())
}
const COLD_LINKS_V15: &str = r#"CREATE TABLE cold.links_v15 (
source_id TEXT NOT NULL,
target_id TEXT NOT NULL,
edge_type TEXT NOT NULL,
valid_from TEXT NOT NULL,
recorded_at TEXT NOT NULL,
valid_to TEXT NOT NULL,
weight REAL NOT NULL CHECK (weight >= 0.0 AND weight < 9e999 AND typeof(weight) = 'real'),
properties TEXT NOT NULL,
branch_id TEXT NOT NULL DEFAULT 'main',
PRIMARY KEY (source_id, target_id, edge_type, valid_from, recorded_at, branch_id)
)"#;
async fn cold_links_keyed_by_lineage(conn: &libsql::Connection) -> Result<bool> {
let mut rows = conn.query("PRAGMA cold.table_info(links)", ()).await?;
while let Some(row) = rows.next().await? {
let named = row.get::<String>(1).is_ok_and(|name| name == "branch_id");
if named && row.get::<i64>(5).is_ok_and(|pk| pk > 0) {
return Ok(true);
}
}
Ok(false)
}
async fn cold_has_branch(conn: &libsql::Connection, table: &str) -> Result<bool> {
let mut rows = conn
.query(&format!("PRAGMA cold.table_info({table})"), ())
.await?;
while let Some(row) = rows.next().await? {
if row.get::<String>(1).is_ok_and(|name| name == "branch_id") {
return Ok(true);
}
}
Ok(false)
}
const ARCHIVED_KEYS: &str = "archived_keys";
async fn collect_archived_keys(
tx: &libsql::Transaction,
clause: &str,
params: impl libsql::params::IntoParams,
) -> Result<()> {
tx.execute(&format!("DROP TABLE IF EXISTS temp.{ARCHIVED_KEYS}"), ())
.await?;
tx.execute(
&format!(
"CREATE TEMP TABLE {ARCHIVED_KEYS} AS \
SELECT DISTINCT {key} FROM links WHERE {clause}",
key = crate::integrity::rebuild::PROJECTION_KEY
),
params,
)
.await?;
Ok(())
}
async fn repair_archived_keys(tx: &libsql::Transaction) -> Result<()> {
crate::integrity::rebuild::repair_keys_within(tx, ARCHIVED_KEYS).await?;
tx.execute(&format!("DROP TABLE temp.{ARCHIVED_KEYS}"), ())
.await?;
Ok(())
}
async fn archive_session(
conn: &libsql::Connection,
cutoff: &str,
archived_at: &str,
) -> Result<ArchiveReport> {
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await?;
for ddl in COLD_SCHEMA {
tx.execute(ddl, ()).await?;
}
upgrade_cold_lineage(&tx).await?;
tx.execute(&format!("CREATE TABLE {ARCHIVE_SESSION_MARKER} (x)"), ())
.await?;
let links_archived = tx
.execute(
&format!(
"INSERT OR IGNORE INTO cold.links
(source_id, target_id, edge_type, valid_from, recorded_at,
valid_to, weight, properties, branch_id)
SELECT source_id, target_id, edge_type, valid_from, recorded_at,
valid_to, weight, properties, branch_id
FROM links WHERE {LINKS_ARCHIVABLE}"
),
libsql::named_params! {":cutoff": cutoff},
)
.await? as usize;
collect_archived_keys(
&tx,
LINKS_ARCHIVABLE,
libsql::named_params! {":cutoff": cutoff},
)
.await?;
let links_deleted = delete_guarded(
&tx,
conn,
&format!("DELETE FROM links WHERE {LINKS_ARCHIVABLE}"),
libsql::named_params! {":cutoff": cutoff},
"links",
)
.await?;
if links_deleted > 0 {
repair_archived_keys(&tx).await?;
}
let concepts_archived = archive_concepts(&tx, conn, cutoff).await?;
let log_entries_archived = tx
.execute(
&format!(
"INSERT OR IGNORE INTO cold.transaction_log
(seq_id, table_name, entity_id, operation, payload, recorded_at, branch_id)
SELECT seq_id, table_name, entity_id, operation, payload, recorded_at, branch_id
FROM transaction_log WHERE {LOG_ARCHIVABLE}"
),
libsql::named_params! {":cutoff": cutoff},
)
.await? as usize;
delete_guarded(
&tx,
conn,
&format!("DELETE FROM transaction_log WHERE {LOG_ARCHIVABLE}"),
libsql::named_params! {":cutoff": cutoff},
"transaction_log",
)
.await?;
let horizon: Option<i64> = tx
.query("SELECT MIN(seq_id) FROM transaction_log", ())
.await?
.next()
.await?
.and_then(|row| row.get(0).ok());
tx.execute(
"INSERT INTO cold.archive_horizon (archived_at, cutoff, horizon) VALUES (?1, ?2, ?3)",
libsql::params![archived_at, cutoff, horizon],
)
.await?;
tx.execute(&format!("DROP TABLE {ARCHIVE_SESSION_MARKER}"), ())
.await?;
tx.commit().await?;
Ok(ArchiveReport {
links_archived,
concepts_archived,
log_entries_archived,
horizon,
})
}
async fn archive_concepts(
tx: &libsql::Transaction,
conn: &libsql::Connection,
cutoff: &str,
) -> Result<usize> {
let moved = tx
.execute(
&format!(
"INSERT OR IGNORE INTO cold.concepts
(rowid_pk, id, title, content, embedding_model,
valid_from, valid_to, recorded_at, retired, branch_id)
SELECT rowid_pk, id, title, content, embedding_model,
valid_from, valid_to, recorded_at, retired, branch_id
FROM concepts WHERE {CONCEPTS_ARCHIVABLE}"
),
libsql::named_params! {":cutoff": cutoff},
)
.await? as usize;
if moved == 0 {
return Ok(0);
}
let mut derived: Vec<String> = vec!["analytics_annotations".to_string()];
let mut rows = tx
.query(
"SELECT name FROM sqlite_master WHERE type = 'table' \
AND name LIKE 'embeddings\\_%' ESCAPE '\\'",
(),
)
.await?;
while let Some(row) = rows.next().await? {
derived.push(row.get::<String>(0)?);
}
drop(rows);
for table in &derived {
tx.execute(
&format!(
"DELETE FROM {table} WHERE concept_id IN \
(SELECT id FROM cold.concepts)"
),
(),
)
.await?;
}
let deleted = delete_guarded(
tx,
conn,
&format!("DELETE FROM concepts WHERE {CONCEPTS_ARCHIVABLE}"),
libsql::named_params! {":cutoff": cutoff},
"concepts",
)
.await? as usize;
debug_assert_eq!(
moved, deleted,
"the predicate selected a different set for the copy than for the delete"
);
Ok(deleted)
}
pub async fn archive_branch(
conn: &libsql::Connection,
branch: &str,
archived_at: &str,
archive_path: &Path,
) -> Result<ArchiveReport> {
crate::temporal::replay::detach_stale_cold(conn).await;
conn.execute(
"ATTACH DATABASE ?1 AS cold",
libsql::params![archive_path.to_string_lossy().as_ref()],
)
.await?;
let result = archive_branch_session(conn, branch, archived_at).await;
if let Err(e) = conn.execute("DETACH DATABASE cold", ()).await {
tracing::warn!("archive_branch: failed to DETACH cold database: {e}");
}
result
}
fn not_archivable(branch: &str, reason: impl Into<String>) -> DbError {
DbError::BranchNotArchivable {
branch: branch.to_string(),
reason: reason.into(),
}
}
async fn any_row(tx: &libsql::Transaction, sql: &str, branch: &str) -> Result<bool> {
Ok(tx
.query(sql, libsql::named_params! {":branch": branch})
.await?
.next()
.await?
.is_some())
}
async fn refuse_unarchivable_branch(tx: &libsql::Transaction, branch: &str) -> Result<()> {
if branch == crate::schema::ddl::MAIN_BRANCH {
return Err(not_archivable(
branch,
"it is the trunk: every lineage's parent chain ends there and every \
default branch_id names it, so there is no ledger left after it goes",
));
}
if !any_row(
tx,
"SELECT 1 FROM branches WHERE branch_id = :branch",
branch,
)
.await?
{
return Err(DbError::UnknownBranch(branch.to_string()));
}
if any_row(
tx,
"SELECT 1 FROM branches WHERE parent_id = :branch",
branch,
)
.await?
{
return Err(not_archivable(
branch,
"it has descendants, which read through it: archiving it would delete \
rows they still believe. Archive the descendants first",
));
}
let mut rows = tx
.query(
"SELECT c.id FROM concepts c
WHERE c.branch_id = :branch
AND EXISTS (
SELECT 1 FROM links l
WHERE l.branch_id <> :branch
AND (l.source_id = c.id OR l.target_id = c.id)
)
LIMIT 1",
libsql::named_params! {":branch": branch},
)
.await?;
if let Some(row) = rows.next().await? {
let id: String = row.get(0)?;
return Err(not_archivable(
branch,
format!(
"concept {id} was minted here and a hot edge on another lineage \
names it. A lineage other lineages still depend on is not \
abandoned; retire those edges first"
),
));
}
Ok(())
}
async fn archive_branch_session(
conn: &libsql::Connection,
branch: &str,
archived_at: &str,
) -> Result<ArchiveReport> {
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await?;
for ddl in COLD_SCHEMA {
tx.execute(ddl, ()).await?;
}
upgrade_cold_lineage(&tx).await?;
refuse_unarchivable_branch(&tx, branch).await?;
tx.execute(&format!("CREATE TABLE {ARCHIVE_SESSION_MARKER} (x)"), ())
.await?;
let links_archived = tx
.execute(
&format!(
"INSERT OR IGNORE INTO cold.links
(source_id, target_id, edge_type, valid_from, recorded_at,
valid_to, weight, properties, branch_id)
SELECT source_id, target_id, edge_type, valid_from, recorded_at,
valid_to, weight, properties, branch_id
FROM links WHERE branch_id = :branch AND branch_id <> '{main}'",
main = crate::schema::ddl::MAIN_BRANCH
),
libsql::named_params! {":branch": branch},
)
.await? as usize;
collect_archived_keys(
&tx,
&format!(
"branch_id = :branch AND branch_id <> '{main}'",
main = crate::schema::ddl::MAIN_BRANCH
),
libsql::named_params! {":branch": branch},
)
.await?;
let links_deleted = delete_guarded(
&tx,
conn,
&format!(
"DELETE FROM links WHERE branch_id = :branch AND branch_id <> '{main}'",
main = crate::schema::ddl::MAIN_BRANCH
),
libsql::named_params! {":branch": branch},
"links",
)
.await?;
if links_deleted > 0 {
repair_archived_keys(&tx).await?;
}
let concepts_archived = archive_branch_concepts(&tx, conn, branch).await?;
let log_entries_archived = tx
.execute(
&format!(
"INSERT OR IGNORE INTO cold.transaction_log
(seq_id, table_name, entity_id, operation, payload, recorded_at, branch_id)
SELECT seq_id, table_name, entity_id, operation, payload, recorded_at, branch_id
FROM transaction_log WHERE branch_id = :branch \
AND branch_id <> '{main}'",
main = crate::schema::ddl::MAIN_BRANCH
),
libsql::named_params! {":branch": branch},
)
.await? as usize;
delete_guarded(
&tx,
conn,
&format!(
"DELETE FROM transaction_log WHERE branch_id = :branch AND branch_id <> '{main}'",
main = crate::schema::ddl::MAIN_BRANCH
),
libsql::named_params! {":branch": branch},
"transaction_log",
)
.await?;
tx.execute(
"INSERT OR IGNORE INTO cold.branches
(branch_id, parent_id, forked_at, created_at, archived_at)
SELECT branch_id, parent_id, forked_at, created_at, ?2
FROM branches WHERE branch_id = ?1",
libsql::params![branch, archived_at],
)
.await?;
delete_guarded(
&tx,
conn,
"DELETE FROM branches WHERE branch_id = :branch",
libsql::named_params! {":branch": branch},
"branches",
)
.await?;
let horizon: Option<i64> = tx
.query("SELECT MIN(seq_id) FROM transaction_log", ())
.await?
.next()
.await?
.and_then(|row| row.get(0).ok());
tx.execute(&format!("DROP TABLE {ARCHIVE_SESSION_MARKER}"), ())
.await?;
tx.commit().await?;
Ok(ArchiveReport {
links_archived,
concepts_archived,
log_entries_archived,
horizon,
})
}
async fn archive_branch_concepts(
tx: &libsql::Transaction,
conn: &libsql::Connection,
branch: &str,
) -> Result<usize> {
let moved = tx
.execute(
"INSERT OR IGNORE INTO cold.concepts
(rowid_pk, id, title, content, embedding_model,
valid_from, valid_to, recorded_at, retired, branch_id)
SELECT rowid_pk, id, title, content, embedding_model,
valid_from, valid_to, recorded_at, retired, branch_id
FROM concepts WHERE branch_id = :branch",
libsql::named_params! {":branch": branch},
)
.await? as usize;
if moved == 0 {
return Ok(0);
}
let mut derived: Vec<String> = vec!["analytics_annotations".to_string()];
let mut rows = tx
.query(
"SELECT name FROM sqlite_master WHERE type = 'table' \
AND name LIKE 'embeddings\\_%' ESCAPE '\\'",
(),
)
.await?;
while let Some(row) = rows.next().await? {
derived.push(row.get::<String>(0)?);
}
drop(rows);
for table in &derived {
tx.execute(
&format!(
"DELETE FROM {table} WHERE concept_id IN \
(SELECT id FROM concepts WHERE branch_id = :branch)"
),
libsql::named_params! {":branch": branch},
)
.await?;
}
let deleted = delete_guarded(
tx,
conn,
"DELETE FROM concepts WHERE branch_id = :branch",
libsql::named_params! {":branch": branch},
"concepts",
)
.await? as usize;
debug_assert_eq!(
moved, deleted,
"the lineage selected a different set for the copy than for the delete"
);
Ok(deleted)
}
#[derive(Clone)]
struct ColdConcept {
old_rowid: i64,
title: String,
content: String,
model: Option<String>,
valid_from: String,
valid_to: String,
recorded_at: String,
retired: i64,
branch_id: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct RehydrateReport {
pub concepts_rehydrated: usize,
pub rowids_reassigned: usize,
}
pub async fn rehydrate(
conn: &libsql::Connection,
ids: &[&str],
archive_path: &Path,
) -> Result<RehydrateReport> {
if ids.is_empty() {
return Ok(RehydrateReport {
concepts_rehydrated: 0,
rowids_reassigned: 0,
});
}
if !archive_present(archive_path) {
return Ok(RehydrateReport {
concepts_rehydrated: 0,
rowids_reassigned: 0,
});
}
crate::temporal::replay::detach_stale_cold(conn).await;
conn.execute(
"ATTACH DATABASE ?1 AS cold",
libsql::params![archive_path.to_string_lossy().as_ref()],
)
.await?;
let result = rehydrate_session(conn, ids).await;
if let Err(e) = conn.execute("DETACH DATABASE cold", ()).await {
tracing::warn!("rehydrate: failed to DETACH cold database: {e}");
}
result
}
async fn rehydrate_session(conn: &libsql::Connection, ids: &[&str]) -> Result<RehydrateReport> {
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await?;
tx.execute(&format!("CREATE TABLE {ARCHIVE_SESSION_MARKER} (x)"), ())
.await?;
let lineage = if cold_has_branch(&tx, "concepts").await? {
"branch_id"
} else {
"'main' AS branch_id"
};
let mut live_lineages = std::collections::HashSet::new();
let mut rows = tx.query("SELECT branch_id FROM branches", ()).await?;
while let Some(row) = rows.next().await? {
live_lineages.insert(row.get::<String>(0)?);
}
drop(rows);
let mut rehydrated = 0usize;
let mut reassigned = 0usize;
for chunk in ids.chunks(HYDRATE_CHUNK) {
let placeholders = std::iter::repeat_n("?", chunk.len())
.collect::<Vec<_>>()
.join(", ");
let bind: Vec<libsql::Value> = chunk.iter().map(|id| libsql::Value::from(*id)).collect();
let mut rows = tx
.query(
&format!(
"SELECT rowid_pk, id, title, content, embedding_model, \
valid_from, valid_to, recorded_at, retired, {lineage} \
FROM cold.concepts WHERE id IN ({placeholders})"
),
libsql::params_from_iter(bind),
)
.await?;
let mut found: std::collections::HashMap<String, ColdConcept> =
std::collections::HashMap::with_capacity(chunk.len());
while let Some(row) = rows.next().await? {
let id: String = row.get(1)?;
found.insert(
id,
ColdConcept {
old_rowid: row.get(0)?,
title: row.get(2)?,
content: row.get(3)?,
model: row.get(4)?,
valid_from: row.get(5)?,
valid_to: row.get(6)?,
recorded_at: row.get(7)?,
retired: row.get(8)?,
branch_id: row.get(9)?,
},
);
}
drop(rows);
let mut moved: Vec<&str> = Vec::with_capacity(found.len());
for id in chunk {
let Some(cold) = found.get(*id) else {
continue;
};
let ColdConcept {
old_rowid,
title,
content,
model,
valid_from,
valid_to,
recorded_at,
retired,
branch_id,
} = cold.clone();
if !live_lineages.contains(&branch_id) {
return Err(DbError::BranchArchived {
branch: branch_id,
concept: (*id).to_string(),
});
}
let taken: i64 = tx
.query(
"SELECT COUNT(*) FROM concepts WHERE rowid_pk = ?1",
libsql::params![old_rowid],
)
.await?
.next()
.await?
.expect("COUNT(*) always returns a row")
.get(0)?;
if taken == 0 {
tx.execute(
"INSERT INTO concepts (rowid_pk, id, title, content, embedding_model, \
valid_from, valid_to, recorded_at, retired, branch_id) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
libsql::params![
old_rowid,
*id,
title.clone(),
content.clone(),
model,
valid_from,
valid_to,
recorded_at,
retired,
branch_id
],
)
.await?;
} else {
tx.execute(
"INSERT INTO concepts (id, title, content, embedding_model, \
valid_from, valid_to, recorded_at, retired, branch_id) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
libsql::params![
*id,
title.clone(),
content.clone(),
model,
valid_from,
valid_to,
recorded_at,
retired,
branch_id
],
)
.await?;
tx.execute(
"INSERT INTO concepts_fts (concepts_fts, rowid, title, content) \
VALUES ('delete', ?1, ?2, ?3)",
libsql::params![old_rowid, title, content],
)
.await?;
reassigned += 1;
}
moved.push(*id);
rehydrated += 1;
}
if !moved.is_empty() {
let placeholders = std::iter::repeat_n("?", moved.len())
.collect::<Vec<_>>()
.join(", ");
let bind: Vec<libsql::Value> =
moved.iter().map(|id| libsql::Value::from(*id)).collect();
tx.execute(
&format!("DELETE FROM cold.concepts WHERE id IN ({placeholders})"),
libsql::params_from_iter(bind),
)
.await?;
}
}
tx.execute(&format!("DROP TABLE {ARCHIVE_SESSION_MARKER}"), ())
.await?;
tx.commit().await?;
Ok(RehydrateReport {
concepts_rehydrated: rehydrated,
rowids_reassigned: reassigned,
})
}
async fn delete_guarded(
tx: &libsql::Transaction,
conn: &libsql::Connection,
sql: &str,
params: impl libsql::params::IntoParams,
table: &str,
) -> Result<u64> {
match tx.execute(sql, params).await {
Ok(n) => Ok(n),
Err(e) => Err(crate::error::classify(conn, e, WriteOp::Delete { table }).await),
}
}
#[cfg(test)]
mod tests {
use super::*;
const EPOCH: &str = "1970-01-01T00:00:00.000000Z";
const OPEN: &str = "9999-12-31T23:59:59.999999Z";
const CLOSED: &str = "1970-01-01T00:30:00.000000Z";
const CUTOFF: &str = "1970-01-01T02:00:00.000000Z";
async fn seeded() -> libsql::Connection {
let db = libsql::Builder::new_local(":memory:")
.build()
.await
.unwrap();
let conn = db.connect().unwrap();
crate::schema::run_migrations(&conn).await.unwrap();
for id in ["a", "b", "c", "e"] {
conn.execute(
"INSERT INTO concepts (id, title, valid_from, recorded_at) \
VALUES (?1, 'n', ?2, ?2)",
libsql::params![id, EPOCH],
)
.await
.unwrap();
}
for (target, valid_to, recorded_at) in [
("b", OPEN, EPOCH),
("b", OPEN, "1970-01-01T01:00:00.000000Z"),
("c", CLOSED, EPOCH),
("e", OPEN, EPOCH),
] {
conn.execute(
"INSERT INTO links (source_id, target_id, edge_type, valid_from, valid_to, \
weight, properties, recorded_at) VALUES ('a', ?1, 'LINKS', ?2, ?3, 1.0, '{}', ?4)",
libsql::params![target, EPOCH, valid_to, recorded_at],
)
.await
.unwrap();
}
conn
}
#[tokio::test]
async fn the_collected_keys_are_only_the_ones_the_delete_disturbs() {
let conn = seeded().await;
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await
.unwrap();
collect_archived_keys(
&tx,
LINKS_ARCHIVABLE,
libsql::named_params! {":cutoff": CUTOFF},
)
.await
.unwrap();
let mut rows = tx
.query(
&format!("SELECT target_id FROM {ARCHIVED_KEYS} ORDER BY target_id"),
(),
)
.await
.unwrap();
let mut targets = Vec::new();
while let Some(row) = rows.next().await.unwrap() {
targets.push(row.get::<String>(0).unwrap());
}
assert_eq!(
targets,
vec!["b".to_string(), "c".to_string()],
"a → e is untouched by this cutoff and the repair has no business \
re-deriving it; a key set this wide is the full rebuild wearing \
the keyed repair's name"
);
}
#[tokio::test]
async fn a_key_asserted_twice_is_collected_once() {
let conn = seeded().await;
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await
.unwrap();
collect_archived_keys(&tx, "1 = 1", ()).await.unwrap();
let n: i64 = tx
.query(&format!("SELECT COUNT(*) FROM {ARCHIVED_KEYS}"), ())
.await
.unwrap()
.next()
.await
.unwrap()
.unwrap()
.get(0)
.unwrap();
assert_eq!(n, 3, "four rows at three keys collected as three keys");
}
}