1#![allow(warnings)]
34
35use std::collections::HashMap;
36use std::path::{Path, PathBuf};
37
38use limnifs_core::codec;
39use limnifs_core::slab_store::SlabStore;
40use limnifs_core::{ContentHandle, ManifestCursor, MetadataBlob};
41
42use crate::sidecar_name;
43
44use crate::config::{ImageMode, WriteConfig};
45use crate::WriteError;
46
47pub struct RwImage {
49 manifest_path: PathBuf,
50 config: WriteConfig,
51 state: Option<OpenState>,
52 inode_map: HashMap<String, u64>,
55 pending_files: HashMap<String, Vec<u8>>,
57 pending_history: Vec<HistoryEntry>,
59 next_inode: u64,
61 base_manifest_hash: Option<[u8; 32]>,
68}
69
70struct OpenState {
73 blob: MetadataBlob,
74 root_inode: u64,
75 slab_store: Option<SlabStore>,
76}
77
78#[derive(Clone, Debug)]
80pub enum HistoryEntry {
81 Add {
82 path: String,
83 inode: u64,
84 size: u64,
85 },
86 Update {
87 path: String,
88 old_inode: u64,
89 new_inode: u64,
90 size: u64,
91 },
92 Delete {
93 path: String,
94 inode: u64,
95 },
96}
97
98impl RwImage {
99 pub fn open(path: &Path, config: WriteConfig) -> Result<Self, WriteError> {
106 cleanup_stale_swap_dir(path);
113
114 let manifest_bytes = std::fs::read(path).map_err(WriteError::Io)?;
115
116 let mut cursor = ManifestCursor::new(&manifest_bytes);
117 let _ = limnifs_core::parse_manifest_header(&mut cursor).map_err(core_to_io)?;
118 let _ = limnifs_core::parse_feature_flags_section(&mut cursor).map_err(core_to_io)?;
119 let meta_ref = limnifs_core::parse_metadata_reference(&mut cursor).map_err(core_to_io)?;
120
121 let blob_bytes: Vec<u8> = if let Some(inline) = meta_ref.inline_metadata.as_ref() {
122 inline.clone()
123 } else {
124 let entry = meta_ref.locators.first().ok_or_else(|| {
125 WriteError::Io(std::io::Error::other(
126 "metadata_reference has neither inline data nor locators",
127 ))
128 })?;
129 let name = sidecar_name(&entry.uri)?;
130 let sidecar = path.parent().unwrap_or_else(|| Path::new(".")).join(name);
131 let wire_bytes = std::fs::read(&sidecar).map_err(WriteError::Io)?;
132 if meta_ref.codec == 0 {
133 wire_bytes
134 } else {
135 codec::decompress(meta_ref.codec, &wire_bytes, meta_ref.uncompressed_len)
136 .map_err(core_to_io)?
137 }
138 };
139
140 let mut blob_cursor = ManifestCursor::new(&blob_bytes);
141 let blob = limnifs_core::parse_metadata_blob(&mut blob_cursor).map_err(core_to_io)?;
142
143 let slab_index = limnifs_core::parse_slab_index(&mut cursor).map_err(core_to_io)?;
144 let slab_store = if slab_index.is_empty() {
145 None
146 } else {
147 Some(SlabStore::load_mmap(path, &slab_index).map_err(core_to_io)?)
148 };
149
150 let root_inode = blob.root_inode_number().ok_or_else(|| {
151 WriteError::Io(std::io::Error::other(
152 "metadata blob: could not identify a unique root directory inode",
153 ))
154 })?;
155
156 let path_index = blob.build_path_index();
157 let next_inode = blob.inodes.iter().map(|i| i.number).max().unwrap_or(0) + 1;
158
159 let mut image = Self {
160 manifest_path: path.to_path_buf(),
161 base_manifest_hash: Some(limnifs_core::hash_section(&manifest_bytes)),
162 config,
163 state: Some(OpenState {
164 blob,
165 root_inode,
166 slab_store,
167 }),
168 inode_map: path_index,
169 pending_files: HashMap::new(),
170 pending_history: Vec::new(),
171 next_inode,
172 };
173 let _ = image.replay_wal_if_present();
175 Ok(image)
176 }
177
178 #[must_use]
181 pub fn create_new(path: &Path, config: WriteConfig) -> Self {
182 Self {
183 manifest_path: path.to_path_buf(),
184 config,
185 state: None,
186 inode_map: HashMap::new(),
187 pending_files: HashMap::new(),
188 pending_history: Vec::new(),
189 next_inode: 1,
190 base_manifest_hash: None,
191 }
192 }
193
194 pub fn add_file(&mut self, path: &str, data: &[u8]) -> Result<u64, WriteError> {
200 let inode = self.next_inode;
201 self.next_inode += 1;
202 let key = normalize_path(path);
203 self.pending_files.insert(key.clone(), data.to_vec());
204 self.inode_map.insert(key.clone(), inode);
205 self.pending_history.push(HistoryEntry::Add {
206 path: key,
207 inode,
208 size: data.len() as u64,
209 });
210 Ok(inode)
211 }
212
213 pub fn update_file(&mut self, path: &str, data: &[u8]) -> Result<(), WriteError> {
219 let key = normalize_path(path);
220 let old_inode = *self.inode_map.get(&key).ok_or_else(|| {
221 WriteError::Io(std::io::Error::other(format!("path not found: {path}")))
222 })?;
223 let new_inode = self.next_inode;
224 self.next_inode += 1;
225 self.pending_files.insert(key.clone(), data.to_vec());
226 self.inode_map.insert(key.clone(), new_inode);
227 self.pending_history.push(HistoryEntry::Update {
228 path: key,
229 old_inode,
230 new_inode,
231 size: data.len() as u64,
232 });
233 Ok(())
234 }
235
236 pub fn delete_file(&mut self, path: &str) -> Result<(), WriteError> {
242 let key = normalize_path(path);
243 let inode = self.inode_map.remove(&key).ok_or_else(|| {
244 WriteError::Io(std::io::Error::other(format!("path not found: {path}")))
245 })?;
246 self.pending_files.remove(&key);
247 self.pending_history
248 .push(HistoryEntry::Delete { path: key, inode });
249 Ok(())
250 }
251
252 pub fn read_file(&self, path: &str) -> Result<Vec<u8>, WriteError> {
259 let state = self.state.as_ref().ok_or_else(|| {
260 WriteError::Io(std::io::Error::other("read_file: image was not opened"))
261 })?;
262 let key = normalize_path(path);
263 if let Some(data) = self.pending_files.get(&key) {
268 return Ok(data.clone());
269 }
270 let inode_num = *self.inode_map.get(&key).ok_or_else(|| {
271 WriteError::Io(std::io::Error::other(format!("path not found: {path}")))
272 })?;
273 let inode = state.blob.inode_by_number(inode_num).ok_or_else(|| {
274 WriteError::Io(std::io::Error::other(format!("inode {inode_num} missing")))
275 })?;
276 match &inode.content_handle {
277 ContentHandle::InlineData(data) => Ok(data.clone()),
278 ContentHandle::SliceMap(slices) => {
279 let store = state.slab_store.as_ref().ok_or_else(|| {
280 WriteError::Io(std::io::Error::other(
281 "read_file: slice-backed file but no slab store",
282 ))
283 })?;
284 let mut out = Vec::new();
285 for slice in slices {
286 let plaintext = store
287 .plaintext_for(slice.drop_id.as_bytes())
288 .ok_or_else(|| {
289 WriteError::Io(std::io::Error::other("drop not in any slab"))
290 })?
291 .map_err(core_to_io)?;
292 out.extend_from_slice(&plaintext);
293 }
294 Ok(out)
295 }
296 _ => Err(WriteError::Io(std::io::Error::other(
297 "read_file: unsupported content handle",
298 ))),
299 }
300 }
301
302 #[must_use]
304 pub fn pending_changes(&self) -> usize {
305 self.pending_history.len()
306 }
307
308 #[must_use]
311 pub fn needs_turnover(&self) -> bool {
312 self.config.turnover_threshold > 0
313 && self.pending_history.len() >= self.config.turnover_threshold as usize
314 }
315
316 #[must_use]
318 pub fn mode(&self) -> &ImageMode {
319 &self.config.mode
320 }
321
322 pub fn commit(&mut self) -> Result<crate::WriteArtifact, WriteError> {
335 self.write_wal()?;
337 let staging = self.staging_dir();
338 self.write_staging_tree(&staging)?;
339 let artifact = crate::write_directory_with_config(&staging, &self.config)?;
340 let _ = std::fs::remove_dir_all(&staging);
341 self.write_artifact(&artifact)?;
342 let _ = std::fs::remove_file(self.wal_path());
346 self.base_manifest_hash = Some(limnifs_core::hash_section(&artifact.bytes));
347 Ok(artifact)
348 }
349
350 pub fn turnover(&self) -> Result<crate::WriteArtifact, WriteError> {
357 let staging = self.staging_dir();
358 if let Some(state) = &self.state {
359 self.write_live_tree_only(state, &staging)?;
360 } else {
361 let _ = std::fs::remove_dir_all(&staging);
362 std::fs::create_dir_all(&staging).map_err(WriteError::Io)?;
363 self.write_pending_only(&staging)?;
364 }
365 let artifact = crate::write_directory_with_config(&staging, &self.config)?;
366 let _ = std::fs::remove_dir_all(&staging);
367 self.write_artifact(&artifact)?;
368 Ok(artifact)
369 }
370
371 fn write_staging_tree(&self, staging: &Path) -> Result<(), WriteError> {
375 let _ = std::fs::remove_dir_all(staging);
376 std::fs::create_dir_all(staging).map_err(WriteError::Io)?;
377
378 if let Some(state) = &self.state {
379 self.write_live_tree(state, staging)?;
380 }
381
382 for entry in &self.pending_history {
385 if let HistoryEntry::Delete { path, .. } = entry {
386 if !self.pending_files.contains_key(path) {
387 let _ = std::fs::remove_file(staging.join(staging_relative(path)));
388 }
389 }
390 }
391
392 for (path, data) in &self.pending_files {
394 let file_path = staging.join(staging_relative(path));
395 if let Some(parent) = file_path.parent() {
396 std::fs::create_dir_all(parent).map_err(WriteError::Io)?;
397 }
398 std::fs::write(&file_path, data).map_err(WriteError::Io)?;
399 }
400 Ok(())
401 }
402
403 fn write_pending_only(&self, staging: &Path) -> Result<(), WriteError> {
404 for (path, data) in &self.pending_files {
405 let file_path = staging.join(staging_relative(path));
406 if let Some(parent) = file_path.parent() {
407 std::fs::create_dir_all(parent).map_err(WriteError::Io)?;
408 }
409 std::fs::write(&file_path, data).map_err(WriteError::Io)?;
410 }
411 Ok(())
412 }
413
414 fn write_live_tree_only(&self, state: &OpenState, staging: &Path) -> Result<(), WriteError> {
417 let _ = std::fs::remove_dir_all(staging);
418 std::fs::create_dir_all(staging).map_err(WriteError::Io)?;
419 self.write_live_tree(state, staging)?;
420 Ok(())
421 }
422
423 fn write_live_tree(&self, state: &OpenState, staging: &Path) -> Result<(), WriteError> {
428 let slab_ref: Option<&dyn limnifs_core::slab_source::SlabSource> = state
429 .slab_store
430 .as_ref()
431 .map(|s| s as &dyn limnifs_core::slab_source::SlabSource);
432 let mut sink = limnifs_core::live_tree::FilesystemSink::new(staging, slab_ref);
433 limnifs_core::live_tree::walk_live_tree(&state.blob, state.root_inode, &mut sink)
434 .map_err(core_to_io)
435 }
436
437 fn write_artifact(&self, artifact: &crate::WriteArtifact) -> Result<(), WriteError> {
452 let parent = self
453 .manifest_path
454 .parent()
455 .unwrap_or_else(|| Path::new("."))
456 .to_path_buf();
457 let staging = parent.join(format!(
463 "{}.new-{}-{}",
464 self.manifest_path
465 .file_name()
466 .and_then(std::ffi::OsStr::to_str)
467 .unwrap_or("image.lim"),
468 std::process::id(),
469 std::time::SystemTime::now()
470 .duration_since(std::time::UNIX_EPOCH)
471 .map(|d| d.as_nanos())
472 .unwrap_or(0),
473 ));
474 let _ = std::fs::remove_dir_all(&staging);
475 std::fs::create_dir_all(&staging).map_err(WriteError::Io)?;
476
477 let manifest_name = self
479 .manifest_path
480 .file_name()
481 .map(std::ffi::OsString::from)
482 .unwrap_or_else(|| std::ffi::OsString::from("image.lim"));
483 let manifest_staging = staging.join(&manifest_name);
484 std::fs::write(&manifest_staging, &artifact.bytes).map_err(WriteError::Io)?;
485
486 let mut slab_names: Vec<std::ffi::OsString> = Vec::new();
487 for slab in &artifact.slabs {
488 let name = sidecar_name(&slab.locator)?;
491 let os_name = std::ffi::OsString::from(name);
492 std::fs::write(staging.join(&os_name), &slab.bytes).map_err(WriteError::Io)?;
493 slab_names.push(os_name);
494 }
495 let sidecar_name: Option<std::ffi::OsString> =
496 if let Some(sidecar) = &artifact.metadata_sidecar {
497 let name = sidecar_name(&sidecar.locator)?;
498 let os_name = std::ffi::OsString::from(name);
499 std::fs::write(staging.join(&os_name), &sidecar.bytes).map_err(WriteError::Io)?;
500 Some(os_name)
501 } else {
502 None
503 };
504
505 if let Some(name) = &sidecar_name {
509 rename_or_fallback(staging.join(name), parent.join(name))?;
510 }
511 for name in &slab_names {
512 rename_or_fallback(staging.join(name), parent.join(name))?;
513 }
514 rename_or_fallback(manifest_staging, self.manifest_path.clone())?;
515
516 let _ = std::fs::remove_dir_all(&staging);
518 Ok(())
519 }
520
521 fn staging_dir(&self) -> PathBuf {
525 let nonce = format!("{}-{}", std::process::id(), self.next_inode);
526 let mut cur = self
527 .manifest_path
528 .parent()
529 .unwrap_or_else(|| Path::new("."))
530 .to_path_buf();
531 loop {
532 if cur.join("Cargo.toml").is_file() {
533 if std::fs::read_to_string(cur.join("Cargo.toml"))
534 .map(|s| s.contains("[workspace]"))
535 .unwrap_or(false)
536 {
537 return cur.join(".scratch").join(format!("limnifs-rw-{nonce}"));
538 }
539 }
540 if !cur.pop() {
541 break;
542 }
543 }
544 self.manifest_path
545 .parent()
546 .unwrap_or_else(|| Path::new("."))
547 .join(".scratch")
548 .join(format!("limnifs-rw-{nonce}"))
549 }
550
551 fn wal_path(&self) -> PathBuf {
553 let parent = self
554 .manifest_path
555 .parent()
556 .unwrap_or_else(|| Path::new("."));
557 let mut name = self
558 .manifest_path
559 .file_name()
560 .map(std::ffi::OsString::from)
561 .unwrap_or_else(|| std::ffi::OsString::from("image.lim"));
562 name.push(".wal");
563 parent.join(name)
564 }
565
566 fn write_wal(&self) -> Result<(), WriteError> {
569 let mut buf: Vec<u8> = Vec::new();
570 buf.extend_from_slice(b"LIMWAL\0\0");
572 buf.extend_from_slice(&self.base_manifest_hash.unwrap_or([0u8; 32]));
573 buf.extend_from_slice(&(self.pending_files.len() as u32).to_le_bytes());
575 for (path, data) in &self.pending_files {
576 write_path_str(&mut buf, path);
577 buf.extend_from_slice(&(data.len() as u64).to_le_bytes());
578 buf.extend_from_slice(data);
579 }
580 buf.extend_from_slice(&(self.pending_history.len() as u32).to_le_bytes());
582 for entry in &self.pending_history {
583 match entry {
584 HistoryEntry::Add { path, .. } => {
585 buf.push(1);
586 write_path_str(&mut buf, path);
587 }
588 HistoryEntry::Update { path, .. } => {
589 buf.push(2);
590 write_path_str(&mut buf, path);
591 }
592 HistoryEntry::Delete { path, .. } => {
593 buf.push(3);
594 write_path_str(&mut buf, path);
595 }
596 }
597 }
598 let wal_tmp = self.wal_path().with_extension("wal.tmp");
600 std::fs::write(&wal_tmp, &buf).map_err(WriteError::Io)?;
601 std::fs::rename(&wal_tmp, self.wal_path()).map_err(WriteError::Io)?;
602 Ok(())
603 }
604
605 fn replay_wal_if_present(&mut self) -> usize {
610 let wal_path = self.wal_path();
611 let Ok(bytes) = std::fs::read(&wal_path) else {
612 return 0;
613 };
614 if bytes.len() < 40 || &bytes[..8] != b"LIMWAL\0\0" {
615 let _ = std::fs::remove_file(&wal_path);
616 return 0;
617 }
618 let mut tag = [0u8; 32];
623 tag.copy_from_slice(&bytes[8..40]);
624 let current = std::fs::read(&self.manifest_path)
625 .map(|b| limnifs_core::hash_section(&b))
626 .ok();
627 if current != Some(tag) {
628 let _ = std::fs::remove_file(&wal_path);
629 return 0;
630 }
631 let mut cursor = WalCursor {
632 bytes: &bytes,
633 pos: 40,
634 };
635 let files_count = match cursor.read_u32_le() {
636 Ok(n) => n as usize,
637 Err(_) => {
638 let _ = std::fs::remove_file(&wal_path);
639 return 0;
640 }
641 };
642 for _ in 0..files_count {
643 let path = match cursor.read_path_str() {
644 Ok(p) => p,
645 Err(_) => break,
646 };
647 let len = match cursor.read_u64_le() {
648 Ok(n) => n as usize,
649 Err(_) => break,
650 };
651 let data = match cursor.read_bytes(len) {
652 Ok(d) => d.to_vec(),
653 Err(_) => break,
654 };
655 self.pending_files.insert(path, data);
656 }
657 let hist_count = match cursor.read_u32_le() {
658 Ok(n) => n as usize,
659 Err(_) => 0,
660 };
661 let mut replayed = 0;
662 for _ in 0..hist_count {
663 let op = match cursor.read_u8() {
664 Ok(b) => b,
665 Err(_) => break,
666 };
667 let path = match cursor.read_path_str() {
668 Ok(p) => p,
669 Err(_) => break,
670 };
671 match op {
672 1 => {
673 let inode = self.next_inode;
674 self.next_inode += 1;
675 let size = self
676 .pending_files
677 .get(&path)
678 .map(|v| v.len() as u64)
679 .unwrap_or(0);
680 self.inode_map.insert(path.clone(), inode);
681 self.pending_history
682 .push(HistoryEntry::Add { path, inode, size });
683 }
684 2 => {
685 let old_inode = self.inode_map.get(&path).copied().unwrap_or(0);
686 let new_inode = self.next_inode;
687 self.next_inode += 1;
688 let size = self
689 .pending_files
690 .get(&path)
691 .map(|v| v.len() as u64)
692 .unwrap_or(0);
693 self.inode_map.insert(path.clone(), new_inode);
694 self.pending_history.push(HistoryEntry::Update {
695 path,
696 old_inode,
697 new_inode,
698 size,
699 });
700 }
701 3 => {
702 let inode = self.inode_map.remove(&path).unwrap_or(0);
703 self.pending_files.remove(&path);
704 self.pending_history
705 .push(HistoryEntry::Delete { path, inode });
706 }
707 _ => break,
708 }
709 replayed += 1;
710 }
711 let _ = std::fs::remove_file(&wal_path);
713 replayed
714 }
715}
716
717fn write_path_str(out: &mut Vec<u8>, s: &str) {
718 let bytes = s.as_bytes();
719 out.extend_from_slice(&(bytes.len() as u32).to_le_bytes());
720 out.extend_from_slice(bytes);
721}
722
723struct WalCursor<'a> {
724 bytes: &'a [u8],
725 pos: usize,
726}
727
728impl<'a> WalCursor<'a> {
729 fn read_u8(&mut self) -> Result<u8, ()> {
730 let b = *self.bytes.get(self.pos).ok_or(())?;
731 self.pos += 1;
732 Ok(b)
733 }
734 fn read_u32_le(&mut self) -> Result<u32, ()> {
735 if self.pos + 4 > self.bytes.len() {
736 return Err(());
737 }
738 let mut arr = [0u8; 4];
739 arr.copy_from_slice(&self.bytes[self.pos..self.pos + 4]);
740 self.pos += 4;
741 Ok(u32::from_le_bytes(arr))
742 }
743 fn read_u64_le(&mut self) -> Result<u64, ()> {
744 if self.pos + 8 > self.bytes.len() {
745 return Err(());
746 }
747 let mut arr = [0u8; 8];
748 arr.copy_from_slice(&self.bytes[self.pos..self.pos + 8]);
749 self.pos += 8;
750 Ok(u64::from_le_bytes(arr))
751 }
752 fn read_bytes(&mut self, len: usize) -> Result<&'a [u8], ()> {
753 if self.pos + len > self.bytes.len() {
754 return Err(());
755 }
756 let slice = &self.bytes[self.pos..self.pos + len];
757 self.pos += len;
758 Ok(slice)
759 }
760 fn read_path_str(&mut self) -> Result<String, ()> {
761 let len = self.read_u32_le()? as usize;
762 let bytes = self.read_bytes(len)?;
763 std::str::from_utf8(bytes).map(String::from).map_err(|_| ())
764 }
765}
766
767fn normalize_path(path: &str) -> String {
770 let trimmed = path.trim_matches('/');
771 if trimmed.is_empty() {
772 "/".to_string()
773 } else {
774 format!("/{trimmed}")
775 }
776}
777
778fn staging_relative(path: &str) -> &str {
781 path.trim_start_matches('/')
782}
783
784fn rename_or_fallback(from: PathBuf, to: PathBuf) -> Result<(), WriteError> {
790 match std::fs::rename(&from, &to) {
791 Ok(()) => Ok(()),
792 Err(e) if e.raw_os_error() == Some(18) => {
793 let bytes = std::fs::read(&from).map_err(WriteError::Io)?;
795 std::fs::write(&to, &bytes).map_err(WriteError::Io)?;
796 let _ = std::fs::remove_file(&from);
797 Ok(())
798 }
799 Err(e) => Err(WriteError::Io(e)),
800 }
801}
802
803fn core_to_io(e: limnifs_core::CoreError) -> WriteError {
804 WriteError::Io(std::io::Error::other(format!("{e}")))
805}
806
807const STALE_SWAP_AGE: std::time::Duration = std::time::Duration::from_secs(30);
821
822fn cleanup_stale_swap_dir(path: &Path) {
823 let Some(name) = path.file_name().and_then(std::ffi::OsStr::to_str) else {
824 return;
825 };
826 let parent = path.parent().unwrap_or_else(|| Path::new("."));
827 let Ok(entries) = std::fs::read_dir(parent) else {
831 return;
832 };
833 for entry in entries.flatten() {
834 let file_name = entry.file_name();
835 let Some(fname) = file_name.to_str() else {
836 continue;
837 };
838 if fname == format!("{name}.new") {
839 let _ = std::fs::remove_dir_all(entry.path());
842 continue;
843 }
844 if let Some(rest) = fname.strip_prefix(&format!("{name}.new-")) {
845 let mtime_old = entry
847 .metadata()
848 .and_then(|m| m.modified())
849 .ok()
850 .and_then(|t| t.elapsed().ok())
851 .is_some_and(|age| age > STALE_SWAP_AGE);
852 if mtime_old {
853 let _ = std::fs::remove_dir_all(entry.path());
854 }
855 let _ = rest;
856 }
857 }
858}
859
860#[cfg(test)]
861mod tests {
862 use super::*;
863 use crate::profile;
864
865 #[test]
866 fn open_cleans_up_stale_new_directory() {
867 let workdir = std::env::temp_dir().join(format!(
871 "limnifs-crash-recovery-{}-{}",
872 std::process::id(),
873 std::time::SystemTime::now()
874 .duration_since(std::time::UNIX_EPOCH)
875 .map(|d| d.as_nanos() as u64)
876 .unwrap_or(0),
877 ));
878 let _ = std::fs::remove_dir_all(&workdir);
879 std::fs::create_dir_all(&workdir).expect("mkdir");
880
881 std::fs::write(workdir.join("data.txt"), b"alpha").expect("src");
883 let manifest = workdir.join("image.lim");
884 let artifact =
885 crate::write_directory_with_config(&workdir, &profile::balanced()).expect("write");
886 std::fs::write(&manifest, &artifact.bytes).expect("manifest");
887 for slab in &artifact.slabs {
888 let name = sidecar_name(&slab.locator).expect("slab locator");
889 std::fs::write(workdir.join(name), &slab.bytes).expect("slab");
890 }
891 if let Some(sidecar) = &artifact.metadata_sidecar {
892 let name = sidecar_name(&sidecar.locator).expect("sidecar locator");
893 std::fs::write(workdir.join(name), &sidecar.bytes).expect("sidecar");
894 }
895
896 let stale = workdir.join("image.lim.new");
898 std::fs::create_dir_all(&stale).expect("mkdir stale");
899 std::fs::write(stale.join("partial.lim"), b"garbage from crashed commit").expect("garbage");
900 assert!(stale.is_dir(), "stale dir exists before open");
901
902 let _image = RwImage::open(&manifest, profile::balanced()).expect("open");
904 assert!(!stale.exists(), "stale dir removed by open");
905
906 let _ = std::fs::remove_dir_all(&workdir);
907 }
908
909 #[test]
910 fn wal_round_trip_recovers_pending_state_after_simulated_crash() {
911 let workdir = std::env::temp_dir().join(format!(
925 "limnifs-wal-rt-{}-{}",
926 std::process::id(),
927 std::time::SystemTime::now()
928 .duration_since(std::time::UNIX_EPOCH)
929 .map(|d| d.as_nanos() as u64)
930 .unwrap_or(0),
931 ));
932 let _ = std::fs::remove_dir_all(&workdir);
933 std::fs::create_dir_all(&workdir).expect("mkdir");
934 std::fs::write(workdir.join("a.txt"), b"alpha").expect("seed a");
935 std::fs::write(workdir.join("b.txt"), b"beta").expect("seed b");
936 let manifest = workdir.join("image.lim");
937 let artifact =
938 crate::write_directory_with_config(&workdir, &profile::balanced()).expect("write");
939 std::fs::write(&manifest, &artifact.bytes).expect("manifest");
940 for slab in &artifact.slabs {
941 let name = sidecar_name(&slab.locator).expect("slab locator");
942 std::fs::write(workdir.join(name), &slab.bytes).expect("slab");
943 }
944 if let Some(sidecar) = &artifact.metadata_sidecar {
945 let name = sidecar_name(&sidecar.locator).expect("sidecar locator");
946 std::fs::write(workdir.join(name), &sidecar.bytes).expect("sidecar");
947 }
948
949 {
951 let mut image = RwImage::open(&manifest, profile::balanced()).expect("open");
952 image.add_file("c.txt", b"gamma").expect("add c");
953 image.update_file("a.txt", b"alpha2").expect("update a");
954 image.delete_file("b.txt").expect("delete b");
955 assert_eq!(image.pending_changes(), 3);
956 image.write_wal().expect("write WAL");
959 assert!(
960 manifest.with_extension("lim.wal").exists() || {
961 let wal = image.wal_path();
963 eprintln!("WAL path: {}", wal.display());
964 wal.exists()
965 }
966 );
967 }
969
970 let image = RwImage::open(&manifest, profile::balanced()).expect("reopen");
972 assert_eq!(
973 image.pending_changes(),
974 3,
975 "WAL should have replayed 3 pending ops"
976 );
977 let _ = std::fs::remove_dir_all(&workdir);
978 }
979
980 #[test]
981 fn rw_image_create_and_add() {
982 let config = profile::balanced();
983 let mut image = RwImage::create_new(Path::new("/tmp/test.lim"), config);
984 let inode = image.add_file("hello.txt", b"hello world").expect("add");
985 assert_eq!(inode, 1);
986 assert_eq!(image.pending_changes(), 1);
987 }
988
989 #[test]
990 fn rw_image_update_and_delete() {
991 let config = profile::balanced();
992 let mut image = RwImage::create_new(Path::new("/tmp/test.lim"), config);
993 image.add_file("file.txt", b"original").expect("add");
994 image.update_file("file.txt", b"updated").expect("update");
995 assert_eq!(image.pending_changes(), 2);
996 image.delete_file("file.txt").expect("delete");
997 assert_eq!(image.pending_changes(), 3);
998 }
999
1000 #[test]
1001 fn rw_image_update_nonexistent_fails() {
1002 let config = profile::balanced();
1003 let mut image = RwImage::create_new(Path::new("/tmp/test.lim"), config);
1004 assert!(image.update_file("nope.txt", b"data").is_err());
1005 }
1006
1007 #[test]
1008 fn rw_image_needs_turnover() {
1009 let mut config = profile::balanced();
1010 config.turnover_threshold = 3;
1011 let mut image = RwImage::create_new(Path::new("/tmp/test.lim"), config);
1012 image.add_file("a", b"data").expect("add");
1013 image.add_file("b", b"data").expect("add");
1014 assert!(!image.needs_turnover());
1015 image.add_file("c", b"data").expect("add");
1016 assert!(image.needs_turnover());
1017 }
1018
1019 fn write_initial(files: &[(&str, &[u8])], config: &WriteConfig) -> PathBuf {
1022 let staging = std::env::temp_dir().join(format!(
1023 "limnifs-rw-test-init-{}-{}",
1024 std::process::id(),
1025 rand_u64()
1026 ));
1027 let _ = std::fs::remove_dir_all(&staging);
1028 std::fs::create_dir_all(&staging).expect("mkdir staging");
1029 for (name, data) in files {
1030 let path = staging.join(name);
1031 if let Some(parent) = path.parent() {
1032 std::fs::create_dir_all(parent).expect("mkdir parent");
1033 }
1034 std::fs::write(&path, data).expect("write file");
1035 }
1036 let manifest = staging.join("image.lim");
1037 let artifact = crate::write_directory_with_config(&staging, config).expect("write");
1038 std::fs::write(&manifest, &artifact.bytes).expect("write manifest");
1039 for slab in &artifact.slabs {
1040 let name = sidecar_name(&slab.locator).expect("locator");
1041 std::fs::write(staging.join(name), &slab.bytes).expect("write slab");
1042 }
1043 if let Some(sidecar) = &artifact.metadata_sidecar {
1044 let name = sidecar_name(&sidecar.locator).expect("locator");
1045 std::fs::write(staging.join(name), &sidecar.bytes).expect("write sidecar");
1046 }
1047 manifest
1048 }
1049
1050 fn rand_u64() -> u64 {
1053 use std::cell::Cell;
1054 use std::time::{SystemTime, UNIX_EPOCH};
1055 thread_local!(static SEED: Cell<u64> = {
1056 let nanos = SystemTime::now()
1057 .duration_since(UNIX_EPOCH)
1058 .map(|d| d.as_nanos() as u64)
1059 .unwrap_or(0);
1060 Cell::new(nanos ^ (std::process::id() as u64).wrapping_mul(0x9E37_79B9_7F4A_7C15))
1061 });
1062 SEED.with(|s| {
1063 let mut x = s.get();
1064 x ^= x << 13;
1065 x ^= x >> 7;
1066 x ^= x << 17;
1067 s.set(x);
1068 x
1069 })
1070 }
1071
1072 #[test]
1073 fn rw_image_open_round_trip() {
1074 let config = profile::balanced();
1075 let manifest = write_initial(
1076 &[("hello.txt", b"hello world"), ("dir/note.txt", b"nested")],
1077 &config,
1078 );
1079 let image = RwImage::open(&manifest, profile::balanced()).expect("open");
1080 assert_eq!(
1081 image.read_file("hello.txt").expect("read hello"),
1082 b"hello world"
1083 );
1084 assert_eq!(
1085 image.read_file("dir/note.txt").expect("read note"),
1086 b"nested"
1087 );
1088 }
1089
1090 #[test]
1091 fn rw_image_commit_adds_file() {
1092 let config = profile::balanced();
1093 let manifest = write_initial(&[("a.txt", b"alpha")], &config);
1094
1095 let mut image = RwImage::open(&manifest, profile::balanced()).expect("open");
1096 image.add_file("b.txt", b"beta").expect("add");
1097 let _ = image.commit().expect("commit");
1098
1099 let reread = RwImage::open(&manifest, profile::balanced()).expect("reopen");
1100 assert_eq!(reread.read_file("a.txt").expect("read a"), b"alpha");
1101 assert_eq!(reread.read_file("b.txt").expect("read b"), b"beta");
1102 }
1103
1104 #[test]
1105 fn concurrent_readers_never_observe_torn_state_during_commit() {
1106 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
1112 use std::sync::Arc;
1113
1114 let config = profile::balanced();
1115 let manifest = write_initial(&[("a.txt", b"gen0")], &config);
1116
1117 let stop = Arc::new(AtomicBool::new(false));
1118 let opens = Arc::new(AtomicUsize::new(0));
1119 let torn = Arc::new(AtomicUsize::new(0));
1120
1121 let mut handles = Vec::new();
1122 for _ in 0..4 {
1123 let m = manifest.clone();
1124 let stop = Arc::clone(&stop);
1125 let opens = Arc::clone(&opens);
1126 let torn = Arc::clone(&torn);
1127 handles.push(std::thread::spawn(move || {
1128 while !stop.load(Ordering::Relaxed) {
1129 match RwImage::open(&m, profile::balanced()) {
1130 Ok(image) => {
1131 opens.fetch_add(1, Ordering::Relaxed);
1132 match image.read_file("a.txt") {
1134 Ok(bytes) if bytes == b"gen0" => {}
1135 Ok(bytes) if bytes == b"genN" => {}
1136 Ok(bytes) => {
1137 eprintln!("torn read: {:?}", String::from_utf8_lossy(&bytes));
1138 torn.fetch_add(1, Ordering::Relaxed);
1139 }
1140 Err(e) => {
1141 eprintln!("torn open->read: {e}");
1142 torn.fetch_add(1, Ordering::Relaxed);
1143 }
1144 }
1145 }
1146 Err(e) => {
1147 eprintln!("torn open: {e}");
1148 torn.fetch_add(1, Ordering::Relaxed);
1149 }
1150 }
1151 }
1152 }));
1153 }
1154
1155 {
1157 let mut image = RwImage::open(&manifest, profile::balanced()).expect("open writer");
1158 image.update_file("a.txt", b"genN").expect("update");
1159 image.commit().expect("commit during readers");
1160 }
1161 std::thread::sleep(std::time::Duration::from_millis(50));
1162 {
1163 let mut image = RwImage::open(&manifest, profile::balanced()).expect("reopen");
1164 image.update_file("a.txt", b"genN").expect("update 2");
1165 image.commit().expect("commit 2");
1166 }
1167
1168 stop.store(true, Ordering::Relaxed);
1169 for h in handles {
1170 h.join().expect("reader thread");
1171 }
1172 assert_eq!(
1173 torn.load(Ordering::Relaxed),
1174 0,
1175 "no reader saw a torn snapshot"
1176 );
1177 assert!(
1178 opens.load(Ordering::Relaxed) > 0,
1179 "readers actually opened the image"
1180 );
1181 }
1182
1183 #[test]
1184 fn rw_image_commit_updates_and_deletes() {
1185 let config = profile::balanced();
1186 let manifest = write_initial(&[("a.txt", b"alpha"), ("b.txt", b"beta")], &config);
1187
1188 let mut image = RwImage::open(&manifest, profile::balanced()).expect("open");
1189 image.update_file("a.txt", b"alpha2").expect("update");
1190 image.delete_file("b.txt").expect("delete");
1191 let _ = image.commit().expect("commit");
1192
1193 let reread = RwImage::open(&manifest, profile::balanced()).expect("reopen");
1194 assert_eq!(reread.read_file("a.txt").expect("read a"), b"alpha2");
1195 assert!(reread.read_file("b.txt").is_err(), "b.txt must be gone");
1196 }
1197
1198 #[test]
1199 fn rw_image_turnover_preserves_tree() {
1200 let config = profile::max_write();
1201 let manifest = write_initial(
1202 &[("hello.txt", b"hello world"), ("dir/note.txt", b"nested")],
1203 &config,
1204 );
1205
1206 let image = RwImage::open(&manifest, profile::max_write()).expect("open");
1207 let _ = image.turnover().expect("turnover");
1208
1209 let reread = RwImage::open(&manifest, profile::max_write()).expect("reopen");
1210 assert_eq!(
1211 reread.read_file("hello.txt").expect("read hello"),
1212 b"hello world"
1213 );
1214 assert_eq!(
1215 reread.read_file("dir/note.txt").expect("read note"),
1216 b"nested"
1217 );
1218 }
1219}