use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, OnceLock};
use crate::config::StoreConfig;
use crate::error::AppError;
use super::{LocalStore, Store};
type Snapshot = Arc<Vec<(PathBuf, Vec<u8>)>>;
fn templates() -> &'static Mutex<HashMap<Vec<String>, Snapshot>> {
static TEMPLATES: OnceLock<Mutex<HashMap<Vec<String>, Snapshot>>> = OnceLock::new();
TEMPLATES.get_or_init(Default::default)
}
pub fn open_with_keyspaces(data_dir: &Path, keyspaces: &[&str]) -> Result<Store, AppError> {
let mut key: Vec<String> = keyspaces.iter().map(|s| s.to_string()).collect();
key.sort();
key.dedup();
match copy_and_open(&snapshot_for(&key)?, data_dir) {
Ok(store) => Ok(store),
Err(_) => {
forget(&key);
std::fs::remove_dir_all(data_dir)?;
copy_and_open(&snapshot_for(&key)?, data_dir)
}
}
}
fn copy_and_open(snapshot: &Snapshot, data_dir: &Path) -> Result<Store, AppError> {
for (rel, bytes) in snapshot.iter() {
let path = data_dir.join(rel);
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
std::fs::write(&path, bytes)?;
}
std::fs::create_dir_all(data_dir)?;
Store::open(&StoreConfig {
data_dir: data_dir.to_path_buf(),
})
}
fn published_dir(key: &[String]) -> PathBuf {
let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
for byte in key.join("\n").bytes() {
hash ^= u64::from(byte);
hash = hash.wrapping_mul(0x0100_0000_01b3);
}
std::env::temp_dir()
.join("vti-store-templates")
.join(format!("{}-{hash:016x}", env!("CARGO_PKG_VERSION")))
}
fn forget(key: &[String]) {
templates()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(key);
let _ = std::fs::remove_dir_all(published_dir(key));
}
fn snapshot_for(key: &[String]) -> Result<Snapshot, AppError> {
let mut templates = templates()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(snapshot) = templates.get(key) {
return Ok(snapshot.clone());
}
let published = published_dir(key);
let snapshot = Arc::new(match load(&published) {
Some(files) => files,
None => {
let parent = published
.parent()
.expect("published template dir has a parent");
std::fs::create_dir_all(parent)?;
let mut lock_path = published.clone().into_os_string();
lock_path.push(".lock");
let lock = std::fs::File::create(lock_path)?;
lock.lock()?;
match load(&published) {
Some(files) => files,
None => build_and_publish(key, &published)?,
}
}
});
templates.insert(key.to_vec(), snapshot.clone());
Ok(snapshot)
}
fn load(published: &Path) -> Option<Vec<(PathBuf, Vec<u8>)>> {
let mut files = Vec::new();
match read_tree(published, Path::new(""), &mut files) {
Ok(()) if !files.is_empty() => Some(files),
_ => None,
}
}
fn build_and_publish(
key: &[String],
published: &Path,
) -> Result<Vec<(PathBuf, Vec<u8>)>, AppError> {
static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let parent = published
.parent()
.expect("published template dir has a parent");
let staging = parent.join(format!(
".staging-{}-{}",
std::process::id(),
SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
));
let built = (|| {
{
let store = LocalStore::open(&StoreConfig {
data_dir: staging.clone(),
})?;
for name in key {
store.keyspace(name)?;
}
}
let mut files = Vec::new();
read_tree(&staging, Path::new(""), &mut files)?;
Ok(files)
})();
if built.is_err() || std::fs::rename(&staging, published).is_err() {
let _ = std::fs::remove_dir_all(&staging);
}
built
}
fn read_tree(root: &Path, rel: &Path, out: &mut Vec<(PathBuf, Vec<u8>)>) -> Result<(), AppError> {
for entry in std::fs::read_dir(root.join(rel))? {
let entry = entry?;
let rel = rel.join(entry.file_name());
if entry.file_type()?.is_dir() {
read_tree(root, &rel, out)?;
} else {
out.push((rel.clone(), std::fs::read(root.join(&rel))?));
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
const SET: &[&str] = &["alpha", "beta", "gamma"];
#[tokio::test]
async fn templated_stores_are_independent_and_reopenable() {
let a = tempfile::tempdir().unwrap();
let b = tempfile::tempdir().unwrap();
{
let store_a = open_with_keyspaces(a.path(), SET).unwrap();
let store_b = open_with_keyspaces(b.path(), SET).unwrap();
store_a
.keyspace("beta")
.unwrap()
.insert("k", &"in a")
.await
.unwrap();
let in_b: Option<String> = store_b.keyspace("beta").unwrap().get("k").await.unwrap();
assert_eq!(in_b, None, "a write to one fixture leaked into another");
store_a.persist().await.unwrap();
}
let reopened = Store::open(&StoreConfig {
data_dir: a.path().to_path_buf(),
})
.unwrap();
let got: Option<String> = reopened.keyspace("beta").unwrap().get("k").await.unwrap();
assert_eq!(got.as_deref(), Some("in a"));
}
#[tokio::test]
async fn keyspaces_outside_the_template_still_open() {
let dir = tempfile::tempdir().unwrap();
let store = open_with_keyspaces(dir.path(), SET).unwrap();
let ks = store.keyspace("delta").unwrap();
ks.insert("k", &1u32).await.unwrap();
assert_eq!(ks.get::<u32>("k").await.unwrap(), Some(1));
}
#[tokio::test]
async fn an_unopenable_published_template_is_rebuilt() {
let set = [format!("unopenable-{}", std::process::id())];
let published = published_dir(&set);
std::fs::create_dir_all(&published).unwrap();
std::fs::write(published.join("version"), b"not a fjall marker").unwrap();
let dir = tempfile::tempdir().unwrap();
let store = open_with_keyspaces(dir.path(), &[set[0].as_str()]).unwrap();
let ks = store.keyspace(&set[0]).unwrap();
ks.insert("k", &1u32).await.unwrap();
assert_eq!(ks.get::<u32>("k").await.unwrap(), Some(1));
let _ = std::fs::remove_dir_all(&published);
}
}