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(())
+}