use std::path::PathBuf;
use genegraph_storage::commit::{
lock_file_for_metadata, with_commit_actor, with_file_lock, with_metadata_file_lock,
};
use genegraph_storage::generations::write_json_atomic;
fn scratch_dir(tag: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!(
"api_public_{tag}_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
dir
}
#[tokio::test]
async fn downstream_metadata_cycle_runs_under_both_locks() {
let dir = scratch_dir("composed");
let metadata_path = dir.join("ds__g1_metadata.json");
write_json_atomic(&metadata_path, r#"{"count":0}"#).unwrap();
assert_eq!(
lock_file_for_metadata(&metadata_path)
.file_name()
.unwrap(),
"ds__g1_metadata.lock"
);
let lock_path = lock_file_for_metadata(&metadata_path);
let (held_tx, held_rx) = std::sync::mpsc::channel::<()>();
let (proceed_tx, proceed_rx) = std::sync::mpsc::channel::<()>();
let lock_for_foreign = lock_path.clone();
let foreign = std::thread::spawn(move || {
tokio::runtime::Builder::new_current_thread()
.build()
.unwrap()
.block_on(async {
with_file_lock(&lock_for_foreign, move || {
held_tx.send(()).unwrap();
proceed_rx.recv().unwrap();
Ok(())
})
.await
.unwrap()
})
});
held_rx.recv_timeout(std::time::Duration::from_secs(5)).unwrap();
let md = metadata_path.clone();
let md_ref = md.clone();
let cycle_task = tokio::spawn(async move {
with_metadata_file_lock(&md_ref, move || async move {
let mut doc: serde_json::Value =
serde_json::from_str(&std::fs::read_to_string(&md).unwrap()).unwrap();
let count = doc["count"].as_u64().unwrap();
doc["count"] = serde_json::json!(count + 1);
write_json_atomic(&md, &doc.to_string()).unwrap();
Ok(())
})
.await
.unwrap();
});
std::thread::sleep(std::time::Duration::from_millis(200));
assert!(
!cycle_task.is_finished(),
"composed cycle must wait for the foreign lock holder"
);
proceed_tx.send(()).unwrap();
foreign.join().unwrap();
cycle_task.await.unwrap();
assert_eq!(
std::fs::read_to_string(&metadata_path).unwrap(),
r#"{"count":1}"#,
"the RMW must execute exactly once, under both locks"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn in_process_actor_cycle_is_publicly_callable() {
let dir = scratch_dir("actor");
let metadata_path = dir.join("ds__g1_metadata.json");
let result: genegraph_storage::StorageResult<()> =
with_commit_actor(&metadata_path, || async { Ok(()) }).await;
result.unwrap();
let _ = std::fs::remove_dir_all(&dir);
}