use crate::{
db::{kv_store::KvStore, sled_store::SledStore, trees::Tree, v1_types::propvals_v1_to_v2},
errors::AtomicResult,
};
pub fn migrate_maybe(store: &SledStore, base_domain: Option<&str>) -> AtomicResult<()> {
const MAX_PASSES: usize = 16;
for pass in 0..MAX_PASSES {
let mut migrated_something = false;
for tree in store.raw_db().tree_names() {
match String::from_utf8_lossy(&tree).as_ref() {
"resources" => {
v0_to_v1(store)?;
migrated_something = true;
}
"reference_index" => {
ref_v0_to_v1(store)?;
migrated_something = true;
}
"resources_v1" => {
resources_v1_to_v2(store)?;
migrated_something = true;
}
"resources_v2" => {
resources_v2_to_v3(store, base_domain)?;
migrated_something = true;
}
"watched_queries" | "members_index" => {
query_index_v1_to_v2(store)?;
migrated_something = true;
}
"watched_queries_v2" | "members_index_v2" => {
query_index_v2_to_v3(store)?;
migrated_something = true;
}
_other => {}
}
}
if !migrated_something {
return Ok(());
}
tracing::info!("Migration pass {} complete, checking for more...", pass + 1);
}
Err(format!(
"Migrations did not converge after {MAX_PASSES} passes — a migration is likely not dropping its source tree."
)
.into())
}
fn resources_v1_to_v2(store: &SledStore) -> AtomicResult<()> {
tracing::warn!("Migrating resources from v1 to v2, this may take a while...");
let old_key = "resources_v1";
let old = store.raw_db().open_tree(old_key)?;
let new_key = "resources_v2";
let new = store.raw_db().open_tree(new_key)?;
new.clear()?;
let mut count = 0;
for item in old.into_iter() {
let (subject, propvals_bin) = item.expect("Unable to convert into interable");
let subject: String =
String::from_utf8(subject.to_vec()).expect("Unable to deserialize subject");
let propvals: crate::db::v1_types::PropValsV1 = bincode1::deserialize(&propvals_bin)
.map_err(|e| format!("Migration Error: Failed to deserialize propvals: {}", e))?;
let new_propvals = propvals_v1_to_v2(propvals);
new.insert(
subject.as_bytes(),
rmp_serde::to_vec(&new_propvals)
.map_err(|e| format!("Migration Error: Failed to encode propvals: {}", e))?,
)?;
count += 1;
}
store.raw_db().drop_tree(old_key).map_err(|e| {
tracing::error!("Migration Error: Failed to drop old tree: {}", e);
e
})?;
tracing::info!("Finished migrating {} resources", count);
tracing::info!("clearing index...");
store.clear_tree(Tree::ValPropSub)?;
store.clear_tree(Tree::PropValSub)?;
store.clear_tree(Tree::QueryMembers)?;
store.clear_tree(Tree::WatchedQueries)?;
tracing::info!("Index cleared. It will be rebuilt on the next query.");
Ok(())
}
fn v0_to_v1(store: &SledStore) -> AtomicResult<()> {
tracing::warn!("Migrating resources schema from v0 to v1...");
let new = store.raw_db().open_tree("resources_v1")?;
let old_key = "resources";
let old = store.raw_db().open_tree(old_key)?;
let mut count = 0;
for item in old.into_iter() {
let (subject, resource_bin) = item.expect("Unable to convert into iterable");
let subject: String =
bincode1::deserialize(&subject).expect("Unable to deserialize subject");
new.insert(subject.as_bytes(), resource_bin)?;
count += 1;
}
let resources_tree_len = store.len(Tree::Resources)?;
assert_eq!(
new.len(),
resources_tree_len,
"Not all resources were migrated."
);
assert!(
store.raw_db().drop_tree(old_key)?,
"Old resources tree not properly removed."
);
tracing::warn!("Finished migration of {} resources", count);
Ok(())
}
fn resources_v2_to_v3(store: &SledStore, base_domain: Option<&str>) -> AtomicResult<()> {
tracing::warn!("Migrating resources from v2 to v3, this may take a while...");
let old_key = "resources_v2";
let old = store.raw_db().open_tree(old_key)?;
let new_key = "resources_v3";
let new = store.raw_db().open_tree(new_key)?;
new.clear()?;
let mut count = 0;
let base_domain = base_domain.unwrap_or("localhost").to_string();
for item in old.into_iter() {
let (subject, propvals_bin) = item.expect("Unable to convert into interable");
let subject_str: String =
String::from_utf8(subject.to_vec()).expect("Unable to deserialize subject");
let new_subject = crate::db::v2_types::string_to_subject(subject_str, &base_domain);
let new_subject_str = new_subject.to_string();
let propvals: crate::db::v2_types::PropValsV2 = rmp_serde::from_slice(&propvals_bin)
.map_err(|e| format!("Migration Error: Failed to deserialize propvals: {}", e))?;
let new_propvals = crate::db::v2_types::propvals_v2_to_v3(propvals, &base_domain);
new.insert(
new_subject_str.as_bytes(),
rmp_serde::to_vec(&new_propvals)
.map_err(|e| format!("Migration Error: Failed to encode propvals: {}", e))?,
)?;
count += 1;
}
store.raw_db().drop_tree(old_key).map_err(|e| {
tracing::error!("Migration Error: Failed to drop old tree: {}", e);
e
})?;
tracing::info!("Finished migrating {} resources", count);
tracing::info!("clearing index...");
store.clear_tree(Tree::ValPropSub)?;
store.clear_tree(Tree::PropValSub)?;
store.clear_tree(Tree::QueryMembers)?;
store.clear_tree(Tree::WatchedQueries)?;
tracing::info!("Index cleared. It will be rebuilt on the next query.");
Ok(())
}
fn query_index_v1_to_v2(store: &SledStore) -> AtomicResult<()> {
tracing::warn!(
"Dropping old query index trees (QueryFilter schema changed — drive field added). \
They will rebuild on next query."
);
let _ = store.raw_db().drop_tree("watched_queries");
let _ = store.raw_db().drop_tree("members_index");
Ok(())
}
fn query_index_v2_to_v3(store: &SledStore) -> AtomicResult<()> {
tracing::warn!(
"Dropping query index v2 trees (QueryFilter key encoding changed — drive is now a key prefix). \
They will rebuild on next query."
);
let _ = store.raw_db().drop_tree("watched_queries_v2");
let _ = store.raw_db().drop_tree("members_index_v2");
Ok(())
}
fn ref_v0_to_v1(store: &SledStore) -> AtomicResult<()> {
tracing::warn!("Rebuilding indexes...");
store.raw_db().drop_tree("reference_index")?;
tracing::warn!("Old reference_index dropped. Index will be rebuilt on next query.");
Ok(())
}