diff --git a/turbopack/crates/turbo-persistence-tools/src/main.rs b/turbopack/crates/turbo-persistence-tools/src/main.rs index 3c9633bd4a69..8a026c99aae5 100644 --- a/turbopack/crates/turbo-persistence-tools/src/main.rs +++ b/turbopack/crates/turbo-persistence-tools/src/main.rs @@ -39,7 +39,7 @@ fn main() -> Result<()> { { println!( " SST {sequence_number:08}.sst: {min_hash:016x} - {max_hash:016x} (p = 1/{})", - u64::MAX / (max_hash - min_hash) + u64::MAX / (max_hash - min_hash + 1) ); println!(" AQMF {aqmf_entries} entries = {} KiB", aqmf_size / 1024); println!( diff --git a/turbopack/crates/turbo-persistence/src/db.rs b/turbopack/crates/turbo-persistence/src/db.rs index 13d426b9971e..e1db68ae0959 100644 --- a/turbopack/crates/turbo-persistence/src/db.rs +++ b/turbopack/crates/turbo-persistence/src/db.rs @@ -525,7 +525,12 @@ impl TurboPersistence { (seq, range.min_hash, range.max_hash, size) }) .collect::>(); - (meta.sequence_number(), meta.family(), ssts) + ( + meta.sequence_number(), + meta.family(), + ssts, + meta.obsolete_sst_files().to_vec(), + ) }) .collect::>(); @@ -606,7 +611,7 @@ impl TurboPersistence { writeln!(log, "Time {time}")?; let span = time.until(Timestamp::now())?; writeln!(log, "Commit {seq:08} {keys_written} keys in {span:#}")?; - for (seq, family, ssts) in new_meta_info { + for (seq, family, ssts, obsolete) in new_meta_info { writeln!(log, "{seq:08} META family:{family}",)?; for (seq, min, max, size) in ssts { writeln!( @@ -615,6 +620,9 @@ impl TurboPersistence { size / 1024 / 1024 )?; } + for seq in obsolete { + writeln!(log, " {seq:08} OBSOLETE SST")?; + } } new_sst_files.sort_unstable_by_key(|(seq, _)| *seq); for (seq, _) in new_sst_files.iter() { diff --git a/turbopack/crates/turbo-persistence/src/meta_file.rs b/turbopack/crates/turbo-persistence/src/meta_file.rs index f48719b97351..9a2d8f7b874b 100644 --- a/turbopack/crates/turbo-persistence/src/meta_file.rs +++ b/turbopack/crates/turbo-persistence/src/meta_file.rs @@ -183,6 +183,8 @@ pub struct MetaFile { family: u32, /// The entries of the file. entries: Vec, + /// The entries that have been marked as obsolete. + obsolete_entries: Vec, /// The obsolete SST files. obsolete_sst_files: Vec, /// The memory mapped file. @@ -247,6 +249,7 @@ impl MetaFile { sequence_number, family, entries, + obsolete_entries: Vec::new(), obsolete_sst_files, mmap, }; @@ -276,11 +279,21 @@ impl MetaFile { pub fn retain_entries(&mut self, mut predicate: impl FnMut(u32) -> bool) -> bool { let old_len = self.entries.len(); - self.entries - .retain(|entry| predicate(entry.sst_data.sequence_number)); + self.entries.retain(|entry| { + if predicate(entry.sst_data.sequence_number) { + true + } else { + self.obsolete_entries.push(entry.sst_data.sequence_number); + false + } + }); old_len != self.entries.len() } + pub fn obsolete_entries(&self) -> &[u32] { + &self.obsolete_entries + } + pub fn has_active_entries(&self) -> bool { !self.entries.is_empty() } diff --git a/turbopack/crates/turbo-persistence/src/meta_file_builder.rs b/turbopack/crates/turbo-persistence/src/meta_file_builder.rs index f0471976f697..eb78f03ff90b 100644 --- a/turbopack/crates/turbo-persistence/src/meta_file_builder.rs +++ b/turbopack/crates/turbo-persistence/src/meta_file_builder.rs @@ -35,17 +35,18 @@ impl MetaFileBuilder { } #[tracing::instrument(level = "trace", skip_all)] - pub fn write(&self, db_path: &Path, seq: u32) -> Result { + pub fn write(self, db_path: &Path, seq: u32) -> Result { let file = db_path.join(format!("{seq:08}.meta")); self.write_internal(&file) .with_context(|| format!("Unable to write meta file {seq:08}.meta")) } - fn write_internal(&self, file: &Path) -> io::Result { + fn write_internal(mut self, file: &Path) -> io::Result { let mut file = BufWriter::new(File::create(file)?); file.write_u32::(0xFE4ADA4A)?; // Magic number file.write_u32::(self.family)?; + self.obsolete_sst_files.sort(); file.write_u32::(self.obsolete_sst_files.len() as u32)?; for obsolete_sst in &self.obsolete_sst_files { file.write_u32::(*obsolete_sst)?; diff --git a/turbopack/crates/turbo-persistence/src/sst_filter.rs b/turbopack/crates/turbo-persistence/src/sst_filter.rs index c0c6520dbaca..ca562880666d 100644 --- a/turbopack/crates/turbo-persistence/src/sst_filter.rs +++ b/turbopack/crates/turbo-persistence/src/sst_filter.rs @@ -19,6 +19,15 @@ impl SstFilter { /// Phase 1: Apply the filter to the meta file and update the state in the filter. pub fn apply_filter(&mut self, meta: &mut MetaFile) { + // Already obsolete entries need to be considered for usage computation + for seq in meta.obsolete_entries() { + if let Some(state) = self.0.get_mut(seq) + && matches!(state, SstState::UnusedObsolete) + { + // the obsolete state is used now + *state = SstState::Obsolete; + } + } meta.retain_entries(|seq| match self.0.entry(seq) { Entry::Occupied(mut e) => { let state = e.get_mut(); diff --git a/turbopack/crates/turbo-persistence/src/tests.rs b/turbopack/crates/turbo-persistence/src/tests.rs index f812eb6ea9f5..8e68e8a30310 100644 --- a/turbopack/crates/turbo-persistence/src/tests.rs +++ b/turbopack/crates/turbo-persistence/src/tests.rs @@ -1,4 +1,4 @@ -use std::time::Instant; +use std::{fs, time::Instant}; use anyhow::Result; use rayon::iter::{IntoParallelIterator, ParallelIterator}; @@ -557,3 +557,99 @@ fn partial_compaction() -> Result<()> { Ok(()) } + +#[test] +fn merge_file_removal() -> Result<()> { + let tempdir = tempfile::tempdir()?; + let path = tempdir.path(); + + let _ = fs::remove_dir_all(path); + + const READ_COUNT: u32 = 2_000; // we'll read every 10th value, so writes are 10x this value + fn put(b: &WriteBatch<(u8, [u8; 4]), 1>, key: u8, value: u32) -> Result<()> { + for i in 0..(READ_COUNT * 10) { + b.put( + 0, + (key, i.to_be_bytes()), + value.to_be_bytes().to_vec().into(), + )?; + } + Ok(()) + } + fn check(db: &TurboPersistence, key: u8, value: u32) -> Result<()> { + for i in 0..READ_COUNT { + // read every 10th item + let i = i * 10; + assert_eq!( + db.get(0, &(key, i.to_be_bytes()))?.as_deref(), + Some(&value.to_be_bytes()[..]), + "Key {key} {i} expected {value}" + ); + } + Ok(()) + } + fn iter_bits(v: u32) -> impl Iterator { + (0..32u8).filter(move |i| v & (1 << i) != 0) + } + + { + println!("--- Init ---"); + let db = TurboPersistence::open(path.to_path_buf())?; + let b = db.write_batch::<_, 1>()?; + for j in 0..=255 { + put(&b, j, 0)?; + } + db.commit_write_batch(b)?; + db.shutdown()?; + } + + let mut expected_values = [0; 256]; + + for i in 1..50 { + println!("--- Iteration {i} ---"); + let i = i * 37; + println!("Add more entries"); + { + let db = TurboPersistence::open(path.to_path_buf())?; + let b = db.write_batch::<_, 1>()?; + for j in iter_bits(i) { + println!("Put {j} = {i}"); + expected_values[j as usize] = i; + put(&b, j, i)?; + } + db.commit_write_batch(b)?; + + for j in 0..32 { + check(&db, j, expected_values[j as usize])?; + } + + db.shutdown()?; + } + + println!("Compaction"); + { + let db = TurboPersistence::open(path.to_path_buf())?; + + db.compact(3.0, 3, u64::MAX)?; + + for j in 0..32 { + check(&db, j, expected_values[j as usize])?; + } + + db.shutdown()?; + } + + println!("Restore check"); + { + let db = TurboPersistence::open(path.to_path_buf())?; + + for j in 0..32 { + check(&db, j, expected_values[j as usize])?; + } + + db.shutdown()?; + } + } + + Ok(()) +}