diff --git a/src/sink.rs b/src/sink.rs index 018b5b5..6316444 100644 --- a/src/sink.rs +++ b/src/sink.rs @@ -213,11 +213,8 @@ pub fn cmd_records_sink( let _ = crate::grain::extend_grain(&dir, &name); } if crate::bark::load(&dir, &name).is_none() { - if let Ok(m) = - crate::bark::with_identity(crate::bark::derived_map(None, &op)) - { - let _ = crate::bark::save(&dir, &name, &m); - } + let _ = + crate::bark::save(&dir, &name, &crate::bark::derived_map(None, &op)); } std::process::exit(code); } @@ -451,7 +448,6 @@ pub fn cmd_records_sink( Value::String(stages.join(" | ")), ); } - let map = crate::bark::with_identity(map)?; crate::bark::save(&dir, &name, &map)?; // Declared retention and declared index, like every writer. (For a @@ -503,3 +499,77 @@ pub fn cmd_records_sink( ); Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::atomic::{AtomicU32, Ordering}; + + static COUNTER: AtomicU32 = AtomicU32::new(0); + + struct TempDir(std::path::PathBuf); + + impl TempDir { + fn new() -> TempDir { + let n = COUNTER.fetch_add(1, Ordering::Relaxed); + let dir = + std::env::temp_dir().join(format!("timberfs-sink-test-{}-{n}", std::process::id())); + fs::create_dir_all(&dir).unwrap(); + TempDir(dir) + } + } + + impl Drop for TempDir { + fn drop(&mut self) { + let _ = fs::remove_dir_all(&self.0); + } + } + + fn a_stream(text: &str) -> Vec { + let mut s = b"\x1estream-start\x1fv=1\x00".to_vec(); + s.extend_from_slice(format!("\x1eentry\x1flen={}\x00{text}\x00", text.len()).as_bytes()); + s.extend_from_slice(b"\x1estream-end\x00"); + s + } + + fn sink_into(dest: &Path, records: &Path) -> anyhow::Result<()> { + let cfg = Config { + chunk_size: 1 << 20, + level: 1, + flush_age_ms: 5_000, + }; + cmd_records_sink( + Some(records), + dest, + cfg, + Delivery::Streaming, + Clock::Now, + false, + None, + None, + "append", + false, + 0.0, + ) + } + + #[test] + fn a_store_with_no_manifest_keeps_the_identity_its_index_already_carries() { + let d = TempDir::new(); + let dest = d.0.join("s.log"); + let records = d.0.join("in.records"); + fs::write(&records, a_stream("hello")).unwrap(); + + sink_into(&dest, &records).unwrap(); + + let (dir, name) = resolve_backing(&dest).unwrap(); + let manifest = crate::bark::load(&dir, &name).expect("the sink writes one"); + let carried = crate::bark::carried_identity(&dir, &name).expect("the index carries one"); + assert_eq!( + manifest.get("id").and_then(|v| v.as_str()), + Some(carried.as_str()), + "one store, one identity" + ); + sink_into(&dest, &records).expect("a second open of the same store"); + } +}