diff --git a/docs/01-app/01-getting-started/08-caching.mdx b/docs/01-app/01-getting-started/08-caching.mdx index cba6ed732ad1..5d277367b963 100644 --- a/docs/01-app/01-getting-started/08-caching.mdx +++ b/docs/01-app/01-getting-started/08-caching.mdx @@ -604,6 +604,6 @@ See [ISR with Cache Components](/docs/app/guides/incremental-static-regeneration ## Bots and crawlers -Browsers receive the static shell instantly. Bots and crawlers are detected by their user agent and handled differently: because they need a complete document, Next.js skips the shell and renders the entire page dynamically at request time, then sends the finished HTML once the render completes. +With [Partial Prerendering](#prerendering), metadata can be dynamic while the rest of the page is prerendered into a static shell. The shell is served immediately and the metadata streams in after. Some bots and crawlers cannot handle this, and need to see the metadata included in the `` tag of the document. Next.js detects them by their user agent, skips the static shell, and renders the page dynamically at request time. -Because the shell is re-rendered instead of reused, work that completed during prerendering now runs at request time for a bot. If part of your shell depends on inputs that only exist while prerendering, such as build-time data or values that are not reachable in the request-time environment, a page that loads for a person can fail to render for a crawler. Make sure the data your shell relies on is also available at request time. See [Bots and crawlers](/docs/app/guides/streaming#bots-and-crawlers) in the Streaming guide for more details. +Data sources available during the build-time prerender have to be available at request time too. Otherwise the render may throw, and these bots receive a 500 error. See [Bots and crawlers](/docs/app/guides/streaming#bots-and-crawlers) in the Streaming guide for which user agents this covers and how to change the list. diff --git a/packages/next/src/shared/lib/errors/canary-only-config-error.test.ts b/packages/next/src/shared/lib/errors/canary-only-config-error.test.ts new file mode 100644 index 000000000000..3ed60f354d2e --- /dev/null +++ b/packages/next/src/shared/lib/errors/canary-only-config-error.test.ts @@ -0,0 +1,40 @@ +import { isStableBuild } from './canary-only-config-error' + +describe('isStableBuild', () => { + afterEach(() => { + delete process.env.__NEXT_VERSION + delete process.env.__NEXT_TEST_MODE + delete process.env.NEXT_PRIVATE_LOCAL_DEV + }) + + it.each([ + // A plain release version. + { version: '16.3.0', expected: true }, + // A numbered preview release published to npm (`next@preview`). + { version: '16.3.0-preview.10', expected: true }, + // A commit preview tarball (scripts/set-preview-version.js) is built from + // an arbitrary canary commit and must not be treated as stable. + { version: '16.4.0-preview-84cee7e6-20260917', expected: false }, + { version: '16.4.0-canary.5', expected: false }, + ])('returns $expected for version $version', ({ version, expected }) => { + process.env.__NEXT_VERSION = version + expect(isStableBuild()).toBe(expected) + }) + + it('returns true when no version is set', () => { + delete process.env.__NEXT_VERSION + expect(isStableBuild()).toBe(true) + }) + + it('returns false in test mode even with a stable version', () => { + process.env.__NEXT_VERSION = '16.3.0' + process.env.__NEXT_TEST_MODE = '1' + expect(isStableBuild()).toBe(false) + }) + + it('returns false for local dev even with a stable version', () => { + process.env.__NEXT_VERSION = '16.3.0' + process.env.NEXT_PRIVATE_LOCAL_DEV = '1' + expect(isStableBuild()).toBe(false) + }) +}) diff --git a/packages/next/src/shared/lib/errors/canary-only-config-error.ts b/packages/next/src/shared/lib/errors/canary-only-config-error.ts index 499fd2b614da..1c28a65da65d 100644 --- a/packages/next/src/shared/lib/errors/canary-only-config-error.ts +++ b/packages/next/src/shared/lib/errors/canary-only-config-error.ts @@ -1,6 +1,12 @@ export function isStableBuild() { + const nextVersion = process.env.__NEXT_VERSION return ( - !process.env.__NEXT_VERSION?.includes('canary') && + !nextVersion?.includes('canary') && + // Commit preview tarballs (e.g. `16.4.0-preview-84cee7e6-20260917`, see + // scripts/set-preview-version.js) are built from arbitrary canary commits, + // so they are not stable. Numbered preview releases published to npm + // (e.g. `16.3.0-preview.10`) are stable. + !nextVersion?.includes('-preview-') && !process.env.__NEXT_TEST_MODE && !process.env.NEXT_PRIVATE_LOCAL_DEV ) diff --git a/turbopack/crates/turbo-persistence/README.md b/turbopack/crates/turbo-persistence/README.md index 79514f3a2f4d..57b033790520 100644 --- a/turbopack/crates/turbo-persistence/README.md +++ b/turbopack/crates/turbo-persistence/README.md @@ -254,10 +254,11 @@ The checksum is verified on the compressed data **before** decompression when th ## Reading -Reading start from the current sequence number and goes downwards. +Opened meta files are stored in per-family shards. A lookup scans only the requested family's meta +files, from newest to oldest; there is no ordering dependency between families. - We have all SST files memory mapped -- for i = CURRENT sequence number .. 0 +- for each meta file of the queried key family, newest first - Check AMQF from SST file for key existence -> if not continue - let block = 0 - loop @@ -325,6 +326,10 @@ Since the process might exit unexpectedly, to avoid "forgetting" to delete the S We limit the number of SST files that are merged at once to avoid long compactions. +Compaction keeps meta files incremental: a new meta file only describes SST files that were merged +or moved. Metadata for untouched SST files stays in its existing meta file. When every active entry +in an old meta file is superseded, that meta file is retired in the same compaction commit. + Full example: Example: diff --git a/turbopack/crates/turbo-persistence/benches/mod.rs b/turbopack/crates/turbo-persistence/benches/mod.rs index 5ec4a29069da..bc4a5e03efaf 100644 --- a/turbopack/crates/turbo-persistence/benches/mod.rs +++ b/turbopack/crates/turbo-persistence/benches/mod.rs @@ -840,6 +840,83 @@ fn bench_read_get_multiple(c: &mut Criterion) { // Compaction Benchmarks // ============================================================================= +fn family_benchmark_key(family: u32, commit: u32, item: u32) -> [u8; 12] { + let mut key = [0; 12]; + key[..4].copy_from_slice(&family.to_be_bytes()); + key[4..8].copy_from_slice(&commit.to_be_bytes()); + key[8..].copy_from_slice(&item.to_be_bytes()); + key +} + +fn bench_family_sharding(c: &mut Criterion) { + const FAMILIES: usize = 4; + const COMMITS: u32 = 100; + let entries_per_commit = scaled(1_000); + let db = LazyLock::new(|| { + let tempdir = tempfile::tempdir().unwrap(); + let config = TpDbConfig { + family_configs: ["family-0", "family-1", "family-2", "family-3"].map(|name| { + FamilyConfig { + name, + kind: FamilyKind::SingleValue, + compression: Compression::Lz4, + } + }), + ..TpDbConfig::new() + }; + let db = TurboPersistence::::open_with_config( + tempdir.path().to_path_buf(), + config, + ) + .unwrap(); + for commit in 0..COMMITS { + for family in 0..FAMILIES as u32 { + let batch = db.write_batch().unwrap(); + for item in 0..entries_per_commit as u32 { + batch + .put( + family, + family_benchmark_key(family, commit, item), + item.to_be_bytes().to_vec().into(), + ) + .unwrap(); + } + db.commit_write_batch(batch).unwrap(); + } + } + (tempdir, db) + }); + + let mut group = c.benchmark_group("read/family_sharding"); + group.measurement_time(Duration::from_secs(5)); + let hit = family_benchmark_key(2, COMMITS - 1, 0); + let miss = family_benchmark_key(2, COMMITS, 0); + let batch_hits = (0..64) + .map(|item| family_benchmark_key(2, COMMITS - 1, item)) + .collect::>(); + let batch_misses = (0..64) + .map(|item| family_benchmark_key(2, COMMITS, item)) + .collect::>(); + + group.bench_function("get/hit", |b| { + let (_, db) = &*db; + b.iter(|| black_box(db.get(2, black_box(&hit)).unwrap())) + }); + group.bench_function("get/miss", |b| { + let (_, db) = &*db; + b.iter(|| black_box(db.get(2, black_box(&miss)).unwrap())) + }); + group.bench_function("batch_get/hit_64", |b| { + let (_, db) = &*db; + b.iter(|| black_box(db.batch_get(2, black_box(&batch_hits)).unwrap())) + }); + group.bench_function("batch_get/miss_64", |b| { + let (_, db) = &*db; + b.iter(|| black_box(db.batch_get(2, black_box(&batch_misses)).unwrap())) + }); + group.finish(); +} + fn bench_compaction(c: &mut Criterion) { let mut group = c.benchmark_group("compaction"); // Compaction is expensive, reduce sample size @@ -1499,6 +1576,6 @@ fn bench_block_cache(c: &mut Criterion) { criterion_group!( name = benches; config = Criterion::default(); - targets = bench_write, bench_write_multi_value, bench_read_get, bench_read_batch_get, bench_read_get_multiple, bench_compaction, bench_compaction_multi_value, bench_qfilter, bench_static_sorted_file_lookup, bench_block_cache + targets = bench_write, bench_write_multi_value, bench_read_get, bench_read_batch_get, bench_read_get_multiple, bench_family_sharding, bench_compaction, bench_compaction_multi_value, bench_qfilter, bench_static_sorted_file_lookup, bench_block_cache ); criterion_main!(benches); diff --git a/turbopack/crates/turbo-persistence/src/db.rs b/turbopack/crates/turbo-persistence/src/db.rs index 4240feae144f..531224fa0769 100644 --- a/turbopack/crates/turbo-persistence/src/db.rs +++ b/turbopack/crates/turbo-persistence/src/db.rs @@ -309,7 +309,7 @@ pub struct TurboPersistence { /// The inner state of the database. Writing will update that. inner: RwLock>, /// A flag to indicate if the database is empty (no meta files). This is an atomic mirror of - /// `inner.meta_files.is_empty()` to avoid taking a lock on the hot path. + /// `inner.is_empty()` to avoid taking a lock on the hot path. is_empty: AtomicBool, /// Tracks whether a write operation is in progress or has permanently failed. /// `None` = idle, `Some(Active)` = in progress, `Some(Error)` = permanently disabled. @@ -334,8 +334,10 @@ pub struct TurboPersistence { /// The inner state of the database. struct Inner { - /// The list of meta files in the database. This is used to derive the SST files. - meta_files: Vec, + /// The list of meta files in the database sharded by family. This is used to derive the SST + /// files. Each family's files are in ascending sequence order; there are no ordering + /// constraints across families. + meta_files_by_family: [Vec; FAMILIES], /// The current sequence number for the database. current_sequence_number: u32, /// The in progress set of hashes of keys that have been accessed. @@ -344,6 +346,43 @@ struct Inner { accessed_key_hashes: [DashSet>; FAMILIES], } +impl Inner { + fn is_empty(&self) -> bool { + self.meta_files_by_family.iter().all(Vec::is_empty) + } + + fn push_meta_file(&mut self, meta_file: MetaFile) { + let family = meta_file.family() as usize; + debug_assert!(family < FAMILIES, "meta file family is out of bounds"); + let shard = &mut self.meta_files_by_family[family]; + debug_assert!( + shard.last().is_none_or(|previous| { + previous.sequence_number() < meta_file.sequence_number() + }), + "meta file appended out of sequence order for family {family}" + ); + shard.push(meta_file); + } + + #[cfg(debug_assertions)] + fn debug_assert_meta_invariants(&self) { + for (family, meta_files) in self.meta_files_by_family.iter().enumerate() { + debug_assert!( + meta_files + .iter() + .all(|meta| meta.family() as usize == family), + "meta file stored in the wrong family shard" + ); + debug_assert!( + meta_files + .windows(2) + .all(|pair| pair[0].sequence_number() < pair[1].sequence_number()), + "meta files in family {family} are not in ascending sequence order" + ); + } + } +} + pub struct CommitOptions { new_meta_files: Vec, new_sst_files: Vec, @@ -442,7 +481,7 @@ impl TurboPersistence path, read_only, inner: RwLock::new(Inner { - meta_files: Vec::new(), + meta_files_by_family: [(); FAMILIES].map(|_| Vec::new()), current_sequence_number: 0, accessed_key_hashes: [(); FAMILIES] .map(|_| DashSet::with_hasher(BuildNoHashHasher::default())), @@ -650,10 +689,13 @@ impl TurboPersistence } let inner = self.inner.get_mut(); - self.is_empty - .store(meta_files.is_empty(), Ordering::Relaxed); - inner.meta_files = meta_files; + for meta_file in meta_files { + inner.push_meta_file(meta_file); + } + #[cfg(debug_assertions)] + inner.debug_assert_meta_invariants(); inner.current_sequence_number = current; + self.is_empty.store(inner.is_empty(), Ordering::Relaxed); Ok(true) } @@ -797,7 +839,7 @@ impl TurboPersistence /// Clears all caches of the database. pub fn clear_cache(&self) { self.clear_block_caches(); - for meta in self.inner.write().meta_files.iter_mut() { + for meta in self.inner.write().meta_files_by_family.iter_mut().flatten() { meta.clear_cache(); } } @@ -816,7 +858,7 @@ impl TurboPersistence /// Prefetches all SST files which are usually lazy loaded. This can be used to reduce latency /// for the first queries after opening the database. pub fn prepare_all_sst_caches(&self) { - for meta in self.inner.write().meta_files.iter_mut() { + for meta in self.inner.write().meta_files_by_family.iter_mut().flatten() { meta.prepare_sst_cache(); } } @@ -1012,8 +1054,8 @@ impl TurboPersistence // in-memory mutations. The MetaFile in-memory optimization // (retain_entries) is deferred to Phase C. let has_delete_file; - let mut meta_seq_numbers_to_delete = Vec::new(); - let entries_to_remove; + let mut meta_seq_numbers_to_delete = [(); FAMILIES].map(|_| Vec::new()); + let mut entries_to_remove = [(); FAMILIES].map(|_| Vec::new()); // Deleted SST bytes: the caller knows each deleted SST's size when it decides to delete it, // so it's carried on `DeletedFile` and summed here (no scan, no stat). stats.bytes_deleted += sst_files_to_delete.iter().map(|f| f.size).sum::(); @@ -1026,17 +1068,17 @@ impl TurboPersistence { let inner = self.inner.read(); - // (A1) Run the SST filter on existing meta files. This only - // updates the SstFilter state — the MetaFile in-memory layout is - // not modified yet (that happens in Phase C via retain_entries). - // Collects the set of SST entry sequence numbers to remove from - // each meta file, keyed by position in `inner.meta_files`. - entries_to_remove = inner - .meta_files - .iter() - .rev() - .map(|meta_file| sst_filter.apply_filter_collect(meta_file)) - .collect::>(); + // (A1) Run the SST filter on existing meta files. This only updates filter state; the + // MetaFile mutation is deferred to Phase C. Each family's removal list is newest-first, + // matching the filter's required recency order. + for (family, meta_files) in inner.meta_files_by_family.iter().enumerate() { + entries_to_remove[family].extend( + meta_files + .iter() + .rev() + .map(|meta_file| sst_filter.apply_filter_collect(meta_file)), + ); + } // (A2) Determine which meta files are fully obsolete by running // `apply_and_get_remove` in newest-first order. Process new metas @@ -1051,11 +1093,15 @@ impl TurboPersistence "newly created meta file should never be a candidate for removal" ); } - for i in (0..inner.meta_files.len()).rev() { - if sst_filter.apply_and_get_remove(&inner.meta_files[i]) { - meta_seq_numbers_to_delete.push(inner.meta_files[i].sequence_number()); - // Deleted meta bytes, read from the `MetaFile`'s mmap length (no stat). - stats.bytes_deleted += inner.meta_files[i].byte_size(); + for (family, meta_files) in inner.meta_files_by_family.iter().enumerate() { + for i in (0..meta_files.len()).rev() { + // Removal lists are newest-first while each shard is oldest-first. + let to_remove = &entries_to_remove[family][meta_files.len() - 1 - i]; + if sst_filter.apply_and_get_remove_after_removing(&meta_files[i], to_remove) { + meta_seq_numbers_to_delete[family].push(meta_files[i].sequence_number()); + // Deleted meta bytes, read from the `MetaFile`'s mmap length (no stat). + stats.bytes_deleted += meta_files[i].byte_size(); + } } } @@ -1064,7 +1110,9 @@ impl TurboPersistence // delete, which consumes one extra sequence number. has_delete_file = !sst_files_to_delete.is_empty() || !blob_seq_numbers_to_delete.is_empty() - || !meta_seq_numbers_to_delete.is_empty(); + || meta_seq_numbers_to_delete + .iter() + .any(|seqs| !seqs.is_empty()); } // Deleted blob bytes. Unlike SST/meta sizes (both already in memory), blob sizes aren't @@ -1088,19 +1136,24 @@ impl TurboPersistence self.parallel_scheduler.block_in_place(|| { if has_delete_file { sst_seq_numbers_to_delete.sort_unstable(); - meta_seq_numbers_to_delete.sort_unstable(); + for seqs in &mut meta_seq_numbers_to_delete { + seqs.sort_unstable(); + } blob_seq_numbers_to_delete.sort_unstable(); // Write *.del file, marking the selected files as to delete let mut buf = Vec::with_capacity( (sst_seq_numbers_to_delete.len() - + meta_seq_numbers_to_delete.len() + + meta_seq_numbers_to_delete + .iter() + .map(Vec::len) + .sum::() + blob_seq_numbers_to_delete.len()) * size_of::(), ); for seq in sst_seq_numbers_to_delete.iter() { buf.write_u32::(*seq)?; } - for seq in meta_seq_numbers_to_delete.iter() { + for seq in meta_seq_numbers_to_delete.iter().flatten() { buf.write_u32::(*seq)?; } for seq in blob_seq_numbers_to_delete.iter() { @@ -1186,12 +1239,9 @@ impl TurboPersistence "SST DELETED", |&seq| seq, )?; - write_seq_numbers( - &mut log, - &meta_seq_numbers_to_delete, - "META DELETED", - |&seq| seq, - )?; + for seqs in &meta_seq_numbers_to_delete { + write_seq_numbers(&mut log, seqs, "META DELETED", |&seq| seq)?; + } anyhow::Ok(()) })() { eprintln!("turbo-persistence: failed to write LOG after commit {seq:08}: {e:#}"); @@ -1208,31 +1258,38 @@ impl TurboPersistence { let mut inner = self.inner.write(); - // Apply the deferred MetaFile mutations from Phase A1. apply_filter - // was called read-only earlier; now we actually move superseded - // entries from active to obsolete inside each MetaFile. - // entries_to_remove was collected in reverse order, so iterate it - // in reverse to match the forward order of inner.meta_files. - for (meta_file, to_remove) in inner - .meta_files - .iter_mut() - .zip(entries_to_remove.into_iter().rev()) + // Apply the deferred removals oldest-first within each family. + for (meta_files, family_removals) in + inner.meta_files_by_family.iter_mut().zip(entries_to_remove) { - if !to_remove.is_empty() { - meta_file.retain_entries(|seq| !to_remove.contains(&seq)); + for (meta_file, to_remove) in + meta_files.iter_mut().zip(family_removals.into_iter().rev()) + { + if !to_remove.is_empty() { + meta_file.retain_entries(|seq| !to_remove.contains(&seq)); + } } } - inner.meta_files.append(&mut new_meta_files); - if !meta_seq_numbers_to_delete.is_empty() { - let to_delete: HashSet = meta_seq_numbers_to_delete.iter().copied().collect(); - inner - .meta_files - .retain(|meta| !to_delete.contains(&meta.sequence_number())); + for meta_file in new_meta_files.drain(..) { + inner.push_meta_file(meta_file); } + for (meta_files, seqs_to_delete) in inner + .meta_files_by_family + .iter_mut() + .zip(&meta_seq_numbers_to_delete) + { + if !seqs_to_delete.is_empty() { + let to_delete: HashSet = seqs_to_delete.iter().copied().collect(); + meta_files.retain(|meta| !to_delete.contains(&meta.sequence_number())); + } + } + #[cfg(debug_assertions)] + inner.debug_assert_meta_invariants(); inner.current_sequence_number = seq; - self.is_empty - .store(inner.meta_files.is_empty(), Ordering::Relaxed); + self.is_empty.store(inner.is_empty(), Ordering::Relaxed); + // The write guard must be released after publishing the matching empty state. + drop(inner); } // Try to delete superseded files immediately. On Linux/macOS this always @@ -1243,7 +1300,9 @@ impl TurboPersistence Self::try_delete_files(&self.path, &sst_seq_numbers_to_delete, "sst") .map(DeferredDeletion::Sst) .chain( - Self::try_delete_files(&self.path, &meta_seq_numbers_to_delete, "meta") + meta_seq_numbers_to_delete + .iter() + .flat_map(|seqs| Self::try_delete_files(&self.path, seqs, "meta")) .map(DeferredDeletion::Meta), ) .chain( @@ -1263,15 +1322,8 @@ impl TurboPersistence writeln!(log, "New database state:")?; writeln!(log, "FAM | META SEQ | SST SEQ FLAGS | RANGE")?; let inner = self.inner.read(); - let families = inner.meta_files.iter().map(|meta| meta.family()).filter({ - let mut set = HashSet::new(); - move |family| set.insert(*family) - }); - for family in families { - for meta in inner.meta_files.iter() { - if meta.family() != family { - continue; - } + for (family, meta_files) in inner.meta_files_by_family.iter().enumerate() { + for meta in meta_files { let meta_seq = meta.sequence_number(); for (entry, range) in meta.entries().iter().zip(meta.hash_ranges()) { let seq = entry.sequence_number(); @@ -1333,7 +1385,7 @@ impl TurboPersistence let inner = self.inner.read(); sequence_number = AtomicU32::new(inner.current_sequence_number); self.compact_internal( - &inner.meta_files, + &inner.meta_files_by_family, &sequence_number, &mut new_meta_files, &mut new_sst_files, @@ -1370,7 +1422,7 @@ impl TurboPersistence /// Internal function to perform a compaction. fn compact_internal( &self, - meta_files: &[MetaFile], + meta_files_by_family: &[Vec; FAMILIES], sequence_number: &AtomicU32, new_meta_files: &mut Vec, new_sst_files: &mut Vec, @@ -1379,13 +1431,13 @@ impl TurboPersistence keys_written: &mut u64, compact_config: &CompactConfig, ) -> Result<()> { - if meta_files.is_empty() { + if meta_files_by_family.iter().all(Vec::is_empty) { return Ok(()); } struct SstWithRange { + /// Index in the current family's `meta_files_by_family` shard. meta_index: usize, - family: u32, index_in_meta: u32, seq: u32, range: StaticSortedFileRange, @@ -1409,31 +1461,35 @@ impl TurboPersistence } } - let ssts_with_ranges = meta_files + let sst_by_family = meta_files_by_family .iter() .enumerate() - .flat_map(|(meta_index, meta)| { - meta.entries() + .map(|(family, meta_files)| { + meta_files .iter() .enumerate() - .map(move |(index_in_meta, entry)| SstWithRange { - meta_index, - family: meta.family(), - index_in_meta: index_in_meta as u32, - seq: entry.sequence_number(), - range: meta.range(index_in_meta as u32), - size: entry.size(), - flags: entry.flags(), + .flat_map(|(meta_index, meta)| { + debug_assert_eq!( + meta.family() as usize, + family, + "meta file stored in the wrong family shard during compaction" + ); + meta.entries() + .iter() + .enumerate() + .map(move |(index_in_meta, entry)| SstWithRange { + meta_index, + index_in_meta: index_in_meta as u32, + seq: entry.sequence_number(), + range: meta.range(index_in_meta as u32), + size: entry.size(), + flags: entry.flags(), + }) }) + .collect::>() }) .collect::>(); - let mut sst_by_family = [(); FAMILIES].map(|_| Vec::new()); - - for sst in ssts_with_ranges { - sst_by_family[sst.family as usize].push(sst); - } - let path = &self.path; let log_mutex = Mutex::new(()); @@ -1466,7 +1522,12 @@ impl TurboPersistence .parallel_map_collect_owned::<_, _, Result>>( merge_jobs, |(family, ssts_with_ranges, merge_jobs)| { + let meta_files = &meta_files_by_family[family]; let family = family as u32; + debug_assert!( + meta_files.iter().all(|meta| meta.family() == family), + "compaction received a meta file from the wrong family shard" + ); if merge_jobs.is_empty() { return Ok(PartialResultPerFamily { @@ -1485,7 +1546,6 @@ impl TurboPersistence let used_key_hashes: Option = { let filters: Vec> = meta_files .iter() - .filter(|m| m.family() == family) .filter_map(|meta_file| { meta_file.deserialize_used_key_hashes_amqf().transpose() }) @@ -1880,6 +1940,7 @@ impl TurboPersistence let guard = log_mutex.lock(); let mut log = self.open_log()?; writeln!(log, "{family:3} | {meta_seq:08} | Compaction:",)?; + for result in merge_result { match result { PartialMergeResult::Merged { @@ -1943,6 +2004,9 @@ impl TurboPersistence for f in sst_files_to_delete.iter() { meta_file_builder.add_obsolete_sst_file(f.seq); } + // Do not copy `used_key_hashes` into the new meta file. Those marks must expire + // as their source meta files are retired; persisting the merged filter here + // would make keys that were used once stay marked as used forever. let new_meta_file = { let _span = tracing::trace_span!("write meta file").entered(); @@ -2063,7 +2127,13 @@ impl TurboPersistence let key_block_cache = self.key_block_cache(); let value_block_cache = self.value_block_cache(); - for meta in inner.meta_files.iter().rev() { + debug_assert!( + inner.meta_files_by_family[family] + .iter() + .all(|meta| meta.family() as usize == family), + "meta file stored in the wrong family shard while querying family {family}" + ); + for meta in inner.meta_files_by_family[family].iter().rev() { match meta.lookup::( family as u32, hash, @@ -2202,7 +2272,13 @@ impl TurboPersistence let inner = self.inner.read(); let key_block_cache = self.key_block_cache(); let value_block_cache = self.value_block_cache(); - for meta in inner.meta_files.iter().rev() { + debug_assert!( + inner.meta_files_by_family[family] + .iter() + .all(|meta| meta.family() as usize == family), + "meta file stored in the wrong family shard while querying family {family}" + ); + for meta in inner.meta_files_by_family[family].iter().rev() { let _result = meta.batch_lookup( family as u32, keys, @@ -2302,8 +2378,13 @@ impl TurboPersistence pub fn statistics(&self) -> Statistics { let inner = self.inner.read(); Statistics { - meta_files: inner.meta_files.len(), - sst_files: inner.meta_files.iter().map(|m| m.entries().len()).sum(), + meta_files: inner.meta_files_by_family.iter().map(Vec::len).sum(), + sst_files: inner + .meta_files_by_family + .iter() + .flatten() + .map(|meta| meta.entries().len()) + .sum(), key_block_cache: CacheStatistics::new(self.key_block_cache()), value_block_cache: CacheStatistics::new(self.value_block_cache()), hits: self.stats.hits_deleted.load(Ordering::Relaxed) @@ -2321,9 +2402,9 @@ impl TurboPersistence Ok(self .inner .read() - .meta_files + .meta_files_by_family .iter() - .rev() + .flat_map(|meta_files| meta_files.iter().rev()) .map(|meta_file| { let entries = meta_file .entries() diff --git a/turbopack/crates/turbo-persistence/src/meta_file.rs b/turbopack/crates/turbo-persistence/src/meta_file.rs index 3ed7376530df..e2948c70cf71 100644 --- a/turbopack/crates/turbo-persistence/src/meta_file.rs +++ b/turbopack/crates/turbo-persistence/src/meta_file.rs @@ -531,10 +531,6 @@ impl MetaFile { &self.obsolete_entries } - pub fn has_active_entries(&self) -> bool { - !self.entries.is_empty() - } - pub fn obsolete_sst_files(&self) -> &[u32] { &self.obsolete_sst_files } diff --git a/turbopack/crates/turbo-persistence/src/sst_filter.rs b/turbopack/crates/turbo-persistence/src/sst_filter.rs index 933960f7a6d5..520e9de7fda0 100644 --- a/turbopack/crates/turbo-persistence/src/sst_filter.rs +++ b/turbopack/crates/turbo-persistence/src/sst_filter.rs @@ -87,6 +87,18 @@ impl SstFilter { /// Updates the filter state for the next meta file. Returns true if the meta file can be /// removed. pub fn apply_and_get_remove(&mut self, meta: &MetaFile) -> bool { + self.apply_and_get_remove_after_removing(meta, &FxHashSet::default()) + } + + /// Like [`apply_and_get_remove`](Self::apply_and_get_remove), but treats entries in + /// `entries_to_remove` as already removed. Commit uses this while its in-memory mutations are + /// deferred until after `CURRENT` is durable, so a newly-subsumed meta file can be retired in + /// the same commit. + pub fn apply_and_get_remove_after_removing( + &mut self, + meta: &MetaFile, + entries_to_remove: &FxHashSet, + ) -> bool { let mut used = false; for seq in meta.obsolete_sst_files() { if let Entry::Occupied(e) = self.0.entry(*seq) { @@ -100,7 +112,11 @@ impl SstFilter { } } - !used && !meta.has_active_entries() + !used + && !meta + .entries() + .iter() + .any(|entry| !entries_to_remove.contains(&entry.sequence_number())) } } diff --git a/turbopack/crates/turbo-persistence/src/tests.rs b/turbopack/crates/turbo-persistence/src/tests.rs index a3f9c6595de4..524116129d6a 100644 --- a/turbopack/crates/turbo-persistence/src/tests.rs +++ b/turbopack/crates/turbo-persistence/src/tests.rs @@ -9,6 +9,7 @@ use crate::{ constants::{MAX_INLINE_VALUE_SIZE, MAX_MEDIUM_VALUE_SIZE, MAX_SMALL_VALUE_SIZE}, db::{CompactConfig, TurboPersistence, read_current_version}, lookup_entry::IterValue, + meta_file::MetaFile, parallel_scheduler::ParallelScheduler, static_sorted_file::{StaticSortedFileIter, StaticSortedFileMetaData}, write_batch::WriteBatch, @@ -2812,3 +2813,99 @@ fn valued_tombstone_rejects_single_value_families() -> Result<()> { db.shutdown()?; Ok(()) } + +#[rstest] +#[case(true)] +#[case(false)] +fn partial_compaction_retires_fully_consumed_meta_files(#[case] mmap: bool) -> Result<()> { + let tempdir = tempfile::tempdir()?; + let path = tempdir.path(); + let access_mode = if mmap { + AccessMode::Mmap + } else { + AccessMode::File + }; + let db = open_db::<1>(path, mmap)?; + + const KEYS: u32 = 2_000; + for generation in 0..4u32 { + let batch = db.write_batch()?; + for key in 0..KEYS { + batch.put( + 0, + key.to_be_bytes().to_vec(), + generation.to_be_bytes().to_vec().into(), + )?; + } + db.commit_write_batch(batch)?; + if generation == 0 { + // Flush this access into the following commit's used-key-hash AMQF. + assert!(db.get(0, &0u32.to_be_bytes())?.is_some()); + } + } + let before_meta_sequences = db + .meta_info()? + .into_iter() + .map(|meta| meta.sequence_number) + .collect::>(); + assert_eq!(before_meta_sequences.len(), 4); + assert!(before_meta_sequences.iter().any(|&seq| { + MetaFile::open(path, seq, None, access_mode) + .unwrap() + .deserialize_used_key_hashes_amqf() + .unwrap() + .is_some() + })); + + let partial = CompactConfig { + min_merge_count: 2, + optimal_merge_count: 2, + max_merge_count: 2, + max_merge_bytes: u64::MAX, + min_merge_duplication_bytes: 0, + optimal_merge_duplication_bytes: 0, + max_merge_segment_count: 1, + }; + assert!(db.compact(&partial)?.is_some()); + let after_partial = db.meta_info()?; + assert_eq!( + after_partial.len(), + 3, + "two fully consumed meta files should retire while two untouched metas remain" + ); + assert_eq!( + after_partial + .iter() + .filter(|meta| before_meta_sequences.contains(&meta.sequence_number)) + .count(), + 2, + "untouched SST metadata should stay in its two existing meta files" + ); + for key in 0..KEYS { + assert_eq!( + &*db.get(0, &key.to_be_bytes())?.unwrap(), + &3u32.to_be_bytes() + ); + } + + db.full_compact()?; + let fully_compacted = db.meta_info()?; + assert_eq!(fully_compacted.len(), 1); + let compacted_meta = + MetaFile::open(path, fully_compacted[0].sequence_number, None, access_mode)?; + assert!( + compacted_meta.deserialize_used_key_hashes_amqf()?.is_none(), + "used-key marks should expire instead of being copied into compaction output" + ); + drop(db); + + let reopened = open_db::<1>(path, mmap)?; + assert_eq!(reopened.meta_info()?.len(), 1); + for key in 0..KEYS { + assert_eq!( + &*reopened.get(0, &key.to_be_bytes())?.unwrap(), + &3u32.to_be_bytes() + ); + } + Ok(()) +}