use std::fs;
use std::path::Path;
use crate::generations::{
GenerationInfo, generation_name, list_artifact_generations, list_generations, logical_name,
parse_generation, write_json_atomic,
};
use crate::lance_storage_graph::LanceStorageGraph;
use crate::traits::backend::StorageBackend;
use super::tmp_dir;
#[test]
fn test_generation_naming_roundtrip() {
assert_eq!(generation_name("ds_ab12", 0), "ds_ab12__g0");
assert_eq!(generation_name("ds_ab12", 17), "ds_ab12__g17");
assert_eq!(logical_name("ds_ab12__g17"), "ds_ab12");
assert_eq!(logical_name("ds_ab12__g0"), "ds_ab12");
assert_eq!(logical_name("ds_ab12"), "ds_ab12", "no suffix = logical");
assert_eq!(parse_generation("ds_ab12__g17"), Some(17));
assert_eq!(parse_generation("ds_ab12"), None);
assert_eq!(parse_generation("ds_ab12__g1x"), None, "non-digit suffix");
assert_eq!(logical_name("ds__g1x"), "ds__g1x");
assert_eq!(logical_name("ds__g1__g2"), "ds__g1");
assert_eq!(parse_generation("ds__g1__g2"), Some(2));
}
#[tokio::test(flavor = "multi_thread")]
async fn test_write_json_atomic_overwrites_completely() {
let dir = tmp_dir("test_write_json_atomic_overwrites_completely").await;
let path = dir.join("ds__g1_metadata.json");
write_json_atomic(&path, r#"{"v": 1}"#).expect("first publish must succeed");
assert_eq!(fs::read_to_string(&path).unwrap(), r#"{"v": 1}"#);
write_json_atomic(&path, r#"{"v": 2, "files": {"rawinput": {}}}"#)
.expect("republish must succeed");
assert_eq!(
fs::read_to_string(&path).unwrap(),
r#"{"v": 2, "files": {"rawinput": {}}}"#
);
let residue: Vec<_> = fs::read_dir(&dir)
.unwrap()
.filter_map(|e| e.ok())
.filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
.collect();
assert!(residue.is_empty(), "tmp files must not leak: {residue:?}");
let _ = fs::remove_dir_all(&dir);
}
#[test]
fn fsync_dir_succeeds_on_real_directory() {
let dir = std::env::temp_dir().join(format!(
"gen_fsync_dir_probe_{}",
uuid::Uuid::new_v4().simple()
));
fs::create_dir_all(&dir).unwrap();
crate::generations::fsync_dir(&dir).expect("fsync on a real directory must succeed");
let _ = fs::remove_dir_all(&dir);
}
#[tokio::test(flavor = "multi_thread")]
async fn write_json_atomic_publishes_into_fresh_parent_dir() {
let dir = tmp_dir("write_json_atomic_fresh_parent").await;
let path = dir.join("nested").join("deeper").join("ds__g1_metadata.json");
write_json_atomic(&path, r#"{"v": 1}"#).expect("publish into fresh parent dirs must succeed");
assert_eq!(fs::read_to_string(&path).unwrap(), r#"{"v": 1}"#);
write_json_atomic(&path, r#"{"v": 2}"#).expect("republish must succeed");
assert_eq!(fs::read_to_string(&path).unwrap(), r#"{"v": 2}"#);
let residue: Vec<_> = fs::read_dir(path.parent().unwrap())
.unwrap()
.filter_map(|e| e.ok())
.filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
.collect();
assert!(residue.is_empty(), "tmp files must not leak: {residue:?}");
let _ = fs::remove_dir_all(&dir);
}
#[test]
fn write_json_atomic_accepts_bare_filename_without_parent() {
let name = format!("bare_{}_metadata.json", uuid::Uuid::new_v4().simple());
let cwd = std::env::current_dir().unwrap();
write_json_atomic(Path::new(&name), r#"{"bare": true}"#)
.expect("bare-filename publish must keep working");
let contents = fs::read_to_string(cwd.join(&name)).unwrap();
assert_eq!(contents, r#"{"bare": true}"#);
let _ = fs::remove_file(cwd.join(&name));
}
#[tokio::test(flavor = "multi_thread")]
async fn test_list_generations_ignores_orphans() {
let dir = tmp_dir("test_list_generations_ignores_orphans").await;
for (generation, rows) in [(1u64, 10usize), (3, 30)] {
write_json_atomic(
&dir.join(format!("ds__g{generation}_metadata.json")),
&format!(r#"{{"nrows": {rows}}}"#),
)
.unwrap();
}
fs::create_dir_all(dir.join("ds__g2_rawinput.lance")).unwrap();
let committed = list_generations(Path::new(&dir), "ds").await.unwrap();
assert_eq!(
committed,
vec![
GenerationInfo {
generation: 1,
metadata_path: dir.join("ds__g1_metadata.json")
},
GenerationInfo {
generation: 3,
metadata_path: dir.join("ds__g3_metadata.json")
},
],
"ascending, committed only"
);
let artifacts = list_artifact_generations(Path::new(&dir), "ds")
.await
.unwrap();
assert_eq!(artifacts, vec![1, 2, 3], "orphans visible to the sweep");
let _ = fs::remove_dir_all(&dir);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_delete_generation_is_prefix_exact() {
let dir = tmp_dir("test_delete_generation_is_prefix_exact").await;
for name in ["ds__g1_rawinput.lance", "ds2__g1_rawinput.lance"] {
fs::create_dir_all(dir.join(name)).unwrap();
}
write_json_atomic(&dir.join("ds__g1_metadata.json"), "{}").unwrap();
write_json_atomic(&dir.join("ds2__g1_metadata.json"), "{}").unwrap();
crate::generations::delete_generation(Path::new(&dir), "ds", 1)
.await
.expect("delete must succeed");
assert!(!dir.join("ds__g1_rawinput.lance").exists());
assert!(!dir.join("ds__g1_metadata.json").exists());
assert!(
dir.join("ds2__g1_rawinput.lance").exists(),
"sibling intact"
);
assert!(dir.join("ds2__g1_metadata.json").exists(), "sibling intact");
fs::create_dir_all(dir.join("ds__g9_rawinput.lance")).unwrap();
crate::generations::delete_generation(Path::new(&dir), "ds", 9)
.await
.expect("orphan sweep must succeed");
assert!(!dir.join("ds__g9_rawinput.lance").exists());
let _ = fs::remove_dir_all(&dir);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_scoped_generation_routes_artifact_paths() {
let dir = tmp_dir("test_scoped_generation_routes_artifact_paths").await;
let storage = LanceStorageGraph::new(dir.to_string_lossy().to_string(), "ds".to_string())
.scoped_generation(3);
assert_eq!(storage.get_name(), "ds__g3");
assert_eq!(
storage.file_path("rawinput"),
dir.join("ds__g3_rawinput.lance")
);
assert_eq!(
storage.metadata_path(),
dir.join("ds__g3_metadata.json"),
"per-generation metadata = per-generation commit pointer"
);
assert_eq!(crate::generations::logical_name(&storage.get_name()), "ds");
let _ = fs::remove_dir_all(&dir);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_scoped_generation_zero_is_the_build_generation() {
let dir = tmp_dir("test_scoped_generation_zero_is_the_build_generation").await;
let storage = LanceStorageGraph::new(dir.to_string_lossy().to_string(), "ds".to_string())
.scoped_generation(0);
assert_eq!(
storage.file_path("lambdas"),
dir.join("ds__g0_lambdas.lance")
);
assert_eq!(storage.metadata_path(), dir.join("ds__g0_metadata.json"));
let _ = fs::remove_dir_all(&dir);
}
use crate::generations::{delete_generation, pin_generation};
fn seed_committed_generation(
base: &Path,
logical: &str,
generation: u64,
) -> (std::path::PathBuf, std::path::PathBuf) {
let md_path = base.join(format!("{logical}__g{generation}_metadata.json"));
fs::write(&md_path, "{}").unwrap();
let artifact = base.join(format!("{logical}__g{generation}_data.lance"));
fs::create_dir_all(&artifact).unwrap();
let inner = artifact.join("part0.lance");
fs::write(&inner, b"payload").unwrap();
(md_path, artifact)
}
#[tokio::test(flavor = "multi_thread")]
async fn pinned_generation_blocks_sweep_until_dropped() {
let base = tmp_dir("gen_pins").await;
let logical = "pin_ds";
let (md_path, artifact) = seed_committed_generation(&base, logical, 1);
let infos = list_generations(&base, logical).await.unwrap();
assert_eq!(infos.len(), 1, "seeded generation is committed");
let guard = pin_generation(&infos[0]).expect("pin committed generation");
let err = delete_generation(&base, logical, 1).await.unwrap_err();
assert!(
matches!(err, crate::StorageError::InvalidState(_)),
"expected InvalidState, got {err:?}"
);
assert!(
md_path.exists(),
"commit pointer must survive a refused sweep"
);
assert!(artifact.exists(), "artifacts must survive a refused sweep");
drop(guard);
delete_generation(&base, logical, 1).await.unwrap();
assert!(!md_path.exists());
assert!(!artifact.exists());
assert!(list_generations(&base, logical).await.unwrap().is_empty());
}
#[tokio::test(flavor = "multi_thread")]
async fn multiple_pins_require_all_readers_dropped() {
let base = tmp_dir("gen_pins_multi").await;
let logical = "pin_ds_multi";
let (md_path, _) = seed_committed_generation(&base, logical, 2);
let infos = list_generations(&base, logical).await.unwrap();
let g1 = pin_generation(&infos[0]).expect("pin 1");
let g2 = pin_generation(&infos[0]).expect("pin 2");
let err = delete_generation(&base, logical, 2).await.unwrap_err();
assert!(matches!(err, crate::StorageError::InvalidState(_)));
drop(g1);
let err = delete_generation(&base, logical, 2).await.unwrap_err();
assert!(
matches!(err, crate::StorageError::InvalidState(_)),
"still pinned by g2"
);
drop(g2);
delete_generation(&base, logical, 2).await.unwrap();
assert!(!md_path.exists());
}
#[tokio::test(flavor = "multi_thread")]
async fn pin_rejects_uncommitted_generation() {
let base = tmp_dir("gen_pins_orphan").await;
let logical = "pin_ds_orphan";
fs::create_dir_all(base.join(format!("{logical}__g3_data.lance"))).unwrap();
let artifact_gens = list_artifact_generations(&base, logical).await.unwrap();
assert_eq!(artifact_gens, vec![3]);
let info = crate::generations::GenerationInfo {
generation: 3,
metadata_path: base.join(format!("{logical}__g3_metadata.json")),
};
let err = pin_generation(&info).unwrap_err();
assert!(
matches!(err, crate::StorageError::Invalid(_)),
"got {err:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn sweep_during_pinned_read_fails_then_succeeds() {
use std::sync::mpsc;
let base = tmp_dir("gen_pins_concurrent").await;
let logical = "pin_ds_race";
let (md_path, artifact) = seed_committed_generation(&base, logical, 7);
let infos = list_generations(&base, logical).await.unwrap();
let info = infos.into_iter().next().unwrap();
let (pinned_tx, pinned_rx) = mpsc::channel::<crate::generations::GenerationGuard>();
let (release_tx, release_rx) = mpsc::channel::<()>();
let reader = std::thread::spawn(move || {
let guard = pin_generation(&info).expect("reader pin");
pinned_tx.send(guard).unwrap();
release_rx.recv().unwrap(); });
let guard = pinned_rx.recv().unwrap();
let err = delete_generation(&base, logical, 7).await.unwrap_err();
assert!(
matches!(err, crate::StorageError::InvalidState(_)),
"got {err:?}"
);
assert!(md_path.exists());
assert!(artifact.exists());
release_tx.send(()).unwrap();
reader.join().unwrap();
drop(guard);
delete_generation(&base, logical, 7).await.unwrap();
assert!(!md_path.exists());
}
#[test]
fn lock_registries_stay_bounded_under_instance_churn() {
for k in 0..2000u32 {
let path = std::env::temp_dir().join(format!("churn_{k}_metadata.json"));
let (a, b) = crate::commit::registry_sizes();
assert!(
(a + b) < 4096,
"registries grew unbounded: commit={a}, dataset={b} after {k} churns"
);
let _ = path;
}
let (a, b) = crate::commit::registry_sizes();
assert!((a + b) < 4096, "final: commit={a}, dataset={b}");
for k in 0..100u32 {
let md = std::env::temp_dir().join(format!("churn2_{k}_metadata.json"));
let dir = std::env::temp_dir().join(format!("churn2_{k}.lance"));
let fut = crate::commit::with_commit_actor(&md, || async { Ok(()) });
tokio::runtime::Builder::new_current_thread()
.build()
.unwrap()
.block_on(fut)
.unwrap();
crate::commit::with_dataset_write_lock(&dir, || Ok(())).unwrap();
}
let (a, b) = crate::commit::registry_sizes();
assert!((a + b) < 4096, "after real churn: commit={a}, dataset={b}");
}
use crate::generations::{SWEEP_GATE_POST_CHECK, SWEEP_GATE_PRE_LOCK, arm_sweep_gate};
#[tokio::test(flavor = "multi_thread")]
async fn pin_during_sweep_removal_window_is_rejected_not_orphaned() {
let base = tmp_dir("gen_race_post_check").await;
let logical = "race_ds";
let (md_path, artifact) = seed_committed_generation(&base, logical, 4);
let infos = list_generations(&base, logical).await.unwrap();
let info = infos.into_iter().next().unwrap();
let (arrived, release) = arm_sweep_gate(SWEEP_GATE_POST_CHECK);
let sweep_base = base.clone();
let sweep = tokio::spawn(async move { delete_generation(&sweep_base, logical, 4).await });
arrived
.recv_timeout(std::time::Duration::from_secs(5))
.expect("sweep parks at the post-check gate");
let pin = std::thread::spawn(move || pin_generation(&info));
release.send(()).unwrap();
sweep
.await
.unwrap()
.expect("sweep completes once the gate opens");
let pin_result = pin.join().unwrap();
assert!(
matches!(pin_result, Err(crate::StorageError::Invalid(_))),
"pin after retirement must fail validation, got {pin_result:?}"
);
assert!(!md_path.exists(), "generation was retired");
assert!(!artifact.exists());
}
#[tokio::test(flavor = "multi_thread")]
async fn pin_registered_before_sweep_lock_wins_the_race() {
let base = tmp_dir("gen_race_pre_lock").await;
let logical = "race_ds2";
let (md_path, artifact) = seed_committed_generation(&base, logical, 5);
let infos = list_generations(&base, logical).await.unwrap();
let info = infos.into_iter().next().unwrap();
let (arrived, release) = arm_sweep_gate(SWEEP_GATE_PRE_LOCK);
let sweep_base = base.clone();
let sweep = tokio::spawn(async move { delete_generation(&sweep_base, logical, 5).await });
arrived
.recv_timeout(std::time::Duration::from_secs(5))
.expect("sweep parks at the pre-lock gate");
let guard = pin_generation(&info).expect("pin wins the race to the lock");
release.send(()).unwrap();
let err = sweep.await.unwrap().unwrap_err();
assert!(
matches!(err, crate::StorageError::InvalidState(_)),
"sweep must refuse a registered pin, got {err:?}"
);
assert!(md_path.exists());
assert!(artifact.exists());
drop(guard);
delete_generation(&base, logical, 5).await.unwrap();
assert!(!md_path.exists());
}