From 11cb227703e72dba1cd94b8c88fe091c1725c351 Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Fri, 31 Jul 2026 20:35:48 +0800 Subject: [PATCH 1/9] fix(overlaybd): validate index bounds before loading --- storage/overlaybd/src/lsmt/file/helper.rs | 48 ++++++++++++-- storage/overlaybd/src/lsmt/file/readonly.rs | 9 +-- storage/overlaybd/src/lsmt/file/readwrite.rs | 3 +- storage/overlaybd/src/lsmt/file/stack.rs | 20 ++---- storage/overlaybd/src/lsmt/file/tests.rs | 68 +++++++++++++++++++- 5 files changed, 119 insertions(+), 29 deletions(-) diff --git a/storage/overlaybd/src/lsmt/file/helper.rs b/storage/overlaybd/src/lsmt/file/helper.rs index 75707543..287b28bc 100644 --- a/storage/overlaybd/src/lsmt/file/helper.rs +++ b/storage/overlaybd/src/lsmt/file/helper.rs @@ -820,10 +820,39 @@ pub(super) async fn verify_ht( ht.verify_magic() && ht.is_trailer() && ht.is_data_file() && ht.is_sealed(), "trailer magic, trailer type, file type or sealedness doesn't match" ); + validate_index_bounds( + ht.index_offset.get(), + ht.index_size.get(), + file_size, + HEADER_SIZE, + )?; Ok(ht) } } +fn validate_index_bounds(offset: u64, count: u64, file_size: u64, trailer_size: u64) -> Result<()> { + let index_limit = file_size + .checked_sub(trailer_size) + .context("index file boundary underflow")?; + let index_size = count + .checked_mul(size_of::() as u64) + .context("index byte length overflow")?; + + ensure!( + offset >= HEADER_SIZE, + "index offset {offset} is before the file header" + ); + ensure!( + offset <= index_limit, + "index offset {offset} is out of bounds for file limit {index_limit}" + ); + ensure!( + index_size <= index_limit - offset, + "index size {count} is out of bounds at offset {offset}" + ); + Ok(()) +} + async fn load_readonly_layer_metadata(file: Arc) -> Result { let file_size = file.size().await?; let header = verify_ht(&file, false, file_size).await?; @@ -874,22 +903,29 @@ pub(super) async fn load_readonly_layers_metadata( async fn load_index( file: &Arc, offset: u64, - count: usize, + count: u64, reset_tag: bool, ) -> Result> { - if count == 0 { + let count_usize = usize::try_from(count).context("index mapping count does not fit usize")?; + let size_bytes = count_usize + .checked_mul(size_of::()) + .context("index allocation size overflow")?; + + let file_size = file.size().await?; + validate_index_bounds(offset, count, file_size, 0)?; + + if count_usize == 0 { return Ok(Vec::new()); } - let size_bytes = count * size_of::(); let mut raw_bytes = vec![0u8; size_bytes]; let read = file.read_at_into(offset, &mut raw_bytes).await?; ensure!(read >= size_bytes, "Index file too short"); - let mut mappings = Vec::with_capacity(count); + let mut mappings = Vec::with_capacity(count_usize); let stride = size_of::(); - for i in 0..count { + for i in 0..count_usize { let start = i * stride; let end = start + stride; let dm = DiskSegmentMapping::read_from_bytes(&raw_bytes[start..end]) @@ -911,7 +947,7 @@ async fn load_index( pub(super) async fn load_index_and_reset_tags( file: &Arc, offset: u64, - count: usize, + count: u64, ) -> Result> { load_index(file, offset, count, true).await } diff --git a/storage/overlaybd/src/lsmt/file/readonly.rs b/storage/overlaybd/src/lsmt/file/readonly.rs index b1c1060e..06fa614f 100644 --- a/storage/overlaybd/src/lsmt/file/readonly.rs +++ b/storage/overlaybd/src/lsmt/file/readonly.rs @@ -31,12 +31,9 @@ impl LSMTReadOnlyFile { let trailer = verify_ht(&file, true, file_size).await?; - let mappings = load_index_and_reset_tags( - &file, - trailer.index_offset.get(), - trailer.index_size.get() as usize, - ) - .await?; + let mappings = + load_index_and_reset_tags(&file, trailer.index_offset.get(), trailer.index_size.get()) + .await?; let index = Arc::new(ReadOnlyIndex::new(mappings)); let uuid = parse_uuid_field(&trailer.uuid).unwrap_or_else(Uuid::nil); diff --git a/storage/overlaybd/src/lsmt/file/readwrite.rs b/storage/overlaybd/src/lsmt/file/readwrite.rs index 05cc08d8..fe78e1f4 100644 --- a/storage/overlaybd/src/lsmt/file/readwrite.rs +++ b/storage/overlaybd/src/lsmt/file/readwrite.rs @@ -118,8 +118,7 @@ impl LSMTFile { } let count = mapping_area_size / stride; - let mappings = - load_index_and_reset_tags(idx_file, HEADER_SIZE, count as usize).await?; + let mappings = load_index_and_reset_tags(idx_file, HEADER_SIZE, count).await?; for m in mappings { mutable_index.insert(m); diff --git a/storage/overlaybd/src/lsmt/file/stack.rs b/storage/overlaybd/src/lsmt/file/stack.rs index c099799d..1eaf36ec 100644 --- a/storage/overlaybd/src/lsmt/file/stack.rs +++ b/storage/overlaybd/src/lsmt/file/stack.rs @@ -79,14 +79,9 @@ async fn merge_readonly_indexes( .enumerate(), ) .map(|(layer_index, (file, metadata))| async move { - let index_size = usize::try_from(metadata.index_size) - .context("readonly layer index size does not fit usize"); - let result = match index_size { - Ok(index_size) => load_index_and_reset_tags(&file, metadata.index_offset, index_size) - .await - .map(ReadOnlyIndex::new), - Err(err) => Err(err), - }; + let result = load_index_and_reset_tags(&file, metadata.index_offset, metadata.index_size) + .await + .map(ReadOnlyIndex::new); (layer_index, result) }) .buffer_unordered(files.len().min(PARALLEL_LOAD_INDEX)) @@ -304,12 +299,9 @@ pub async fn stack_files( pub async fn open_file_index(file: Arc) -> Result { let file_size = file.size().await?; let trailer = verify_ht(&file, true, file_size).await?; - let mappings = load_index_and_reset_tags( - &file, - trailer.index_offset.get(), - trailer.index_size.get() as usize, - ) - .await?; + let mappings = + load_index_and_reset_tags(&file, trailer.index_offset.get(), trailer.index_size.get()) + .await?; Ok(ReadOnlyIndex::new(mappings)) } diff --git a/storage/overlaybd/src/lsmt/file/tests.rs b/storage/overlaybd/src/lsmt/file/tests.rs index 907ccd0f..bfe87a6b 100644 --- a/storage/overlaybd/src/lsmt/file/tests.rs +++ b/storage/overlaybd/src/lsmt/file/tests.rs @@ -446,6 +446,72 @@ async fn create_sealed_layer( data } +async fn update_trailer(file: &Arc, update: impl FnOnce(&mut HeaderTrailer)) { + let file_size = file.size().await.unwrap(); + let trailer_offset = file_size - HEADER_SIZE; + let bytes = file + .read_at(trailer_offset, size_of::()) + .await + .unwrap(); + let mut trailer = HeaderTrailer::read_from_bytes(&bytes).unwrap(); + update(&mut trailer); + file.write_at(trailer_offset, trailer.as_bytes()) + .await + .unwrap(); +} + +#[tokio::test] +async fn rejects_sealed_index_offset_out_of_bounds() { + let temp_dir = TempDir::new().unwrap(); + let data = create_sealed_layer(&temp_dir, "bad-index-offset", 8192, &[(0, 0x11)]).await; + let file_size = data.size().await.unwrap(); + + update_trailer(&data, |trailer| { + trailer.index_offset = U64::new(file_size); + }) + .await; + + let err = open_file_ro(data as Arc) + .await + .err() + .expect("malformed index offset should be rejected"); + assert_err_contains(&err, "index offset"); +} + +#[tokio::test] +async fn rejects_sealed_index_size_out_of_bounds() { + let temp_dir = TempDir::new().unwrap(); + let data = create_sealed_layer(&temp_dir, "bad-index-size", 8192, &[(0, 0x22)]).await; + let file_size = data.size().await.unwrap(); + let bad_index_size = file_size / size_of::() as u64 + 1; + + update_trailer(&data, |trailer| { + trailer.index_size = U64::new(bad_index_size); + }) + .await; + + let err = open_file_ro(data as Arc) + .await + .err() + .expect("malformed index size should be rejected"); + assert_err_contains(&err, "index size"); +} + +#[tokio::test] +async fn rejects_index_count_overflow_before_allocation() { + let temp_dir = TempDir::new().unwrap(); + let data = create_sealed_layer(&temp_dir, "overflow-index", 8192, &[]).await; + let count = u64::MAX / size_of::() as u64 + 1; + let err = load_index_and_reset_tags(&(data as Arc), HEADER_SIZE, count) + .await + .expect_err("overflowing index count should be rejected"); + let message = err.to_string(); + assert!( + message.contains("overflow") || message.contains("does not fit usize"), + "expected an overflow or conversion error, got {message}" + ); +} + async fn open_sparse_lsmt_env( data_file: Arc, lower_layers: Vec>, @@ -471,7 +537,7 @@ async fn load_base_index(data_file: Arc) -> ReadOnlyIndex { let mappings = load_index_and_reset_tags( &(data_file as Arc), trailer.index_offset.get(), - trailer.index_size.get() as usize, + trailer.index_size.get(), ) .await .expect("Failed to load mappings"); From 12c639850a5c2c07fe1775e8761bb0f972104e1f Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Fri, 31 Jul 2026 20:35:48 +0800 Subject: [PATCH 2/9] test(overlaybd): cover combined index bounds --- storage/overlaybd/src/lsmt/file/tests.rs | 24 +++++++++++++++++++++++- 1 file changed, 23 insertions(+), 1 deletion(-) diff --git a/storage/overlaybd/src/lsmt/file/tests.rs b/storage/overlaybd/src/lsmt/file/tests.rs index bfe87a6b..933ff314 100644 --- a/storage/overlaybd/src/lsmt/file/tests.rs +++ b/storage/overlaybd/src/lsmt/file/tests.rs @@ -478,12 +478,34 @@ async fn rejects_sealed_index_offset_out_of_bounds() { assert_err_contains(&err, "index offset"); } +#[tokio::test] +async fn rejects_sealed_index_offset_before_header() { + let temp_dir = TempDir::new().unwrap(); + let data = create_sealed_layer(&temp_dir, "index-before-header", 8192, &[(0, 0x12)]).await; + + update_trailer(&data, |trailer| { + trailer.index_offset = U64::new(HEADER_SIZE - 1); + }) + .await; + + let err = open_file_ro(data as Arc) + .await + .err() + .expect("index offset inside the header should be rejected"); + assert_err_contains(&err, "index offset"); +} + #[tokio::test] async fn rejects_sealed_index_size_out_of_bounds() { let temp_dir = TempDir::new().unwrap(); let data = create_sealed_layer(&temp_dir, "bad-index-size", 8192, &[(0, 0x22)]).await; let file_size = data.size().await.unwrap(); - let bad_index_size = file_size / size_of::() as u64 + 1; + let trailer = verify_ht(&(data.clone() as Arc), true, file_size) + .await + .unwrap(); + let stride = size_of::() as u64; + let remaining = file_size - HEADER_SIZE - trailer.index_offset.get(); + let bad_index_size = remaining / stride + 1; update_trailer(&data, |trailer| { trailer.index_size = U64::new(bad_index_size); From 9f09aaf1c0bba3f33b29a303da41f9941cea1579 Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Fri, 31 Jul 2026 21:45:00 +0800 Subject: [PATCH 3/9] fix(overlaybd): separate trailer and index validation --- storage/overlaybd/src/lsmt/file/helper.rs | 19 ++++++++++++------- storage/overlaybd/src/lsmt/file/readonly.rs | 6 ++++++ storage/overlaybd/src/lsmt/file/stack.rs | 6 ++++++ storage/overlaybd/src/lsmt/file/tests.rs | 6 ++++++ 4 files changed, 30 insertions(+), 7 deletions(-) diff --git a/storage/overlaybd/src/lsmt/file/helper.rs b/storage/overlaybd/src/lsmt/file/helper.rs index 287b28bc..6ef2e8e1 100644 --- a/storage/overlaybd/src/lsmt/file/helper.rs +++ b/storage/overlaybd/src/lsmt/file/helper.rs @@ -820,17 +820,16 @@ pub(super) async fn verify_ht( ht.verify_magic() && ht.is_trailer() && ht.is_data_file() && ht.is_sealed(), "trailer magic, trailer type, file type or sealedness doesn't match" ); - validate_index_bounds( - ht.index_offset.get(), - ht.index_size.get(), - file_size, - HEADER_SIZE, - )?; Ok(ht) } } -fn validate_index_bounds(offset: u64, count: u64, file_size: u64, trailer_size: u64) -> Result<()> { +pub(super) fn validate_index_bounds( + offset: u64, + count: u64, + file_size: u64, + trailer_size: u64, +) -> Result<()> { let index_limit = file_size .checked_sub(trailer_size) .context("index file boundary underflow")?; @@ -857,6 +856,12 @@ async fn load_readonly_layer_metadata(file: Arc) -> Result) -> Result { let file_size = file.size().await?; let trailer = verify_ht(&file, true, file_size).await?; + validate_index_bounds( + trailer.index_offset.get(), + trailer.index_size.get(), + file_size, + HEADER_SIZE, + )?; let mappings = load_index_and_reset_tags(&file, trailer.index_offset.get(), trailer.index_size.get()) .await?; diff --git a/storage/overlaybd/src/lsmt/file/tests.rs b/storage/overlaybd/src/lsmt/file/tests.rs index 933ff314..66081253 100644 --- a/storage/overlaybd/src/lsmt/file/tests.rs +++ b/storage/overlaybd/src/lsmt/file/tests.rs @@ -471,6 +471,12 @@ async fn rejects_sealed_index_offset_out_of_bounds() { }) .await; + let rw_err = LSMTFile::open(data.clone(), None, None, Vec::new()) + .await + .err() + .expect("malformed sealed layer must not open as writable"); + assert_err_contains(&rw_err, "Sealed"); + let err = open_file_ro(data as Arc) .await .err() From bf11ab00d1c3d7b2420362bc64dbfae506d555a9 Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Fri, 31 Jul 2026 22:27:41 +0800 Subject: [PATCH 4/9] fix(overlaybd): bound index allocations --- storage/overlaybd/src/lsmt/file/helper.rs | 17 ++++++++++++++--- storage/overlaybd/src/lsmt/file/tests.rs | 19 ++++++++++++------- storage/overlaybd/src/lsmt/file/types.rs | 1 + 3 files changed, 27 insertions(+), 10 deletions(-) diff --git a/storage/overlaybd/src/lsmt/file/helper.rs b/storage/overlaybd/src/lsmt/file/helper.rs index 6ef2e8e1..6897c695 100644 --- a/storage/overlaybd/src/lsmt/file/helper.rs +++ b/storage/overlaybd/src/lsmt/file/helper.rs @@ -847,7 +847,7 @@ pub(super) fn validate_index_bounds( ); ensure!( index_size <= index_limit - offset, - "index size {count} is out of bounds at offset {offset}" + "index byte length {index_size} for {count} mappings is out of bounds at offset {offset}" ); Ok(()) } @@ -915,6 +915,10 @@ async fn load_index( let size_bytes = count_usize .checked_mul(size_of::()) .context("index allocation size overflow")?; + ensure!( + size_bytes <= MAX_INDEX_BYTES, + "index byte length {size_bytes} exceeds limit {MAX_INDEX_BYTES}" + ); let file_size = file.size().await?; validate_index_bounds(offset, count, file_size, 0)?; @@ -923,11 +927,18 @@ async fn load_index( return Ok(Vec::new()); } - let mut raw_bytes = vec![0u8; size_bytes]; + let mut raw_bytes = Vec::new(); + raw_bytes + .try_reserve_exact(size_bytes) + .context("failed to reserve index byte buffer")?; + raw_bytes.resize(size_bytes, 0); let read = file.read_at_into(offset, &mut raw_bytes).await?; ensure!(read >= size_bytes, "Index file too short"); - let mut mappings = Vec::with_capacity(count_usize); + let mut mappings = Vec::new(); + mappings + .try_reserve_exact(count_usize) + .context("failed to reserve decoded index mappings")?; let stride = size_of::(); for i in 0..count_usize { diff --git a/storage/overlaybd/src/lsmt/file/tests.rs b/storage/overlaybd/src/lsmt/file/tests.rs index 66081253..6d24f8bb 100644 --- a/storage/overlaybd/src/lsmt/file/tests.rs +++ b/storage/overlaybd/src/lsmt/file/tests.rs @@ -471,12 +471,6 @@ async fn rejects_sealed_index_offset_out_of_bounds() { }) .await; - let rw_err = LSMTFile::open(data.clone(), None, None, Vec::new()) - .await - .err() - .expect("malformed sealed layer must not open as writable"); - assert_err_contains(&rw_err, "Sealed"); - let err = open_file_ro(data as Arc) .await .err() @@ -522,7 +516,7 @@ async fn rejects_sealed_index_size_out_of_bounds() { .await .err() .expect("malformed index size should be rejected"); - assert_err_contains(&err, "index size"); + assert_err_contains(&err, "index byte length"); } #[tokio::test] @@ -540,6 +534,17 @@ async fn rejects_index_count_overflow_before_allocation() { ); } +#[tokio::test] +async fn rejects_index_larger_than_allocation_limit() { + let temp_dir = TempDir::new().unwrap(); + let data = create_sealed_layer(&temp_dir, "oversized-index", 8192, &[]).await; + let count = MAX_INDEX_BYTES as u64 / size_of::() as u64 + 1; + let err = load_index_and_reset_tags(&(data as Arc), HEADER_SIZE, count) + .await + .expect_err("oversized index should be rejected before allocation"); + assert_err_contains(&err, "exceeds limit"); +} + async fn open_sparse_lsmt_env( data_file: Arc, lower_layers: Vec>, diff --git a/storage/overlaybd/src/lsmt/file/types.rs b/storage/overlaybd/src/lsmt/file/types.rs index 634f23b2..968e75b2 100644 --- a/storage/overlaybd/src/lsmt/file/types.rs +++ b/storage/overlaybd/src/lsmt/file/types.rs @@ -12,6 +12,7 @@ pub(super) const ALIGNMENT: u64 = 512; pub(super) const ALIGNMENT_USIZE: usize = 512; pub(super) const HEADER_SIZE: u64 = 4096; pub(super) const MAX_IO_SIZE: usize = 4 * 1024 * 1024; // 4MB +pub(super) const MAX_INDEX_BYTES: usize = 1024 * 1024 * 1024; // 1 GiB pub(super) const ALIGNMENT_4K: usize = 4096; pub(super) const INVALID_SEGMENT_OFFSET: u64 = (1 << 50) - 1; pub(super) const COMPACT_ZERO_DETECTION_ENABLED: bool = false; From 97843b9a31d83a8fa9f8001cba79e88b84fce4ae Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Fri, 31 Jul 2026 23:11:00 +0800 Subject: [PATCH 5/9] fix(overlaybd): budget decoded index memory --- storage/overlaybd/src/lsmt/file/helper.rs | 37 +++++++++++++------- storage/overlaybd/src/lsmt/file/readwrite.rs | 1 + storage/overlaybd/src/lsmt/file/stack.rs | 9 +++++ storage/overlaybd/src/lsmt/file/tests.rs | 10 +++--- storage/overlaybd/src/lsmt/file/types.rs | 3 +- 5 files changed, 43 insertions(+), 17 deletions(-) diff --git a/storage/overlaybd/src/lsmt/file/helper.rs b/storage/overlaybd/src/lsmt/file/helper.rs index 6897c695..9cf9e62f 100644 --- a/storage/overlaybd/src/lsmt/file/helper.rs +++ b/storage/overlaybd/src/lsmt/file/helper.rs @@ -852,6 +852,22 @@ pub(super) fn validate_index_bounds( Ok(()) } +pub(super) fn index_memory_bytes(count: u64) -> Result { + let count = usize::try_from(count).context("index mapping count does not fit usize")?; + count + .checked_mul(size_of::() + size_of::()) + .context("index memory size overflow") +} + +pub(super) fn validate_index_memory(count: u64) -> Result<()> { + let memory_bytes = index_memory_bytes(count)?; + ensure!( + memory_bytes <= MAX_INDEX_MEMORY_BYTES, + "index memory {memory_bytes} exceeds per-layer limit {MAX_INDEX_MEMORY_BYTES}" + ); + Ok(()) +} + async fn load_readonly_layer_metadata(file: Arc) -> Result { let file_size = file.size().await?; let header = verify_ht(&file, false, file_size).await?; @@ -912,14 +928,10 @@ async fn load_index( reset_tag: bool, ) -> Result> { let count_usize = usize::try_from(count).context("index mapping count does not fit usize")?; + validate_index_memory(count)?; let size_bytes = count_usize .checked_mul(size_of::()) .context("index allocation size overflow")?; - ensure!( - size_bytes <= MAX_INDEX_BYTES, - "index byte length {size_bytes} exceeds limit {MAX_INDEX_BYTES}" - ); - let file_size = file.size().await?; validate_index_bounds(offset, count, file_size, 0)?; @@ -1395,12 +1407,6 @@ pub async fn compact_to( compress_raw_index(&mut compact_index); let index_offset = dest_moffset * ALIGNMENT; - let mut index_bytes = Vec::with_capacity(compact_index.len() * size_of::()); - for m in &compact_index { - let dm = DiskSegmentMapping::from_memory(m); - index_bytes.extend_from_slice(dm.as_bytes()); - } - let n_per_block = 4096 / size_of::(); let remainder = compact_index.len() % n_per_block; let padding_count = if remainder > 0 { @@ -1408,12 +1414,19 @@ pub async fn compact_to( } else { 0 }; + let index_size = compact_index.len() as u64 + padding_count as u64; + validate_index_memory(index_size)?; + + let mut index_bytes = Vec::with_capacity(compact_index.len() * size_of::()); + for m in &compact_index { + let dm = DiskSegmentMapping::from_memory(m); + index_bytes.extend_from_slice(dm.as_bytes()); + } append_invalid_index_padding(&mut index_bytes, padding_count); writer.write_all_at(&index_bytes, index_offset).await?; - let index_size = compact_index.len() as u64 + padding_count as u64; let mut trailer_offset = index_offset + index_bytes.len() as u64; if !trailer_offset.is_multiple_of(4096) { diff --git a/storage/overlaybd/src/lsmt/file/readwrite.rs b/storage/overlaybd/src/lsmt/file/readwrite.rs index fe78e1f4..e9d620dd 100644 --- a/storage/overlaybd/src/lsmt/file/readwrite.rs +++ b/storage/overlaybd/src/lsmt/file/readwrite.rs @@ -889,6 +889,7 @@ impl LSMTFile { // HEADER_SIZE + virtual_size; punching the last virtual block ends // exactly at this boundary, so index bytes cannot overlap virtual data. let index_offset = data_file_size; + validate_index_memory(compact_index.len() as u64)?; let mut index_bytes = Vec::with_capacity(compact_index.len() * size_of::()); diff --git a/storage/overlaybd/src/lsmt/file/stack.rs b/storage/overlaybd/src/lsmt/file/stack.rs index 8530cfef..e7a0a2f1 100644 --- a/storage/overlaybd/src/lsmt/file/stack.rs +++ b/storage/overlaybd/src/lsmt/file/stack.rs @@ -61,6 +61,15 @@ async fn merge_readonly_indexes( files.len() == metadata.len(), "readonly files/metadata length mismatch" ); + let total_index_memory = metadata.iter().try_fold(0usize, |total, layer| { + total + .checked_add(index_memory_bytes(layer.index_size)?) + .context("stack index memory size overflow") + })?; + ensure!( + total_index_memory <= MAX_STACK_INDEX_MEMORY_BYTES, + "stack index memory {total_index_memory} exceeds limit {MAX_STACK_INDEX_MEMORY_BYTES}" + ); // `buffer_unordered` returns a stream whose `Future` implementation the // compiler cannot prove is `Send`, even though every input is `Send`. diff --git a/storage/overlaybd/src/lsmt/file/tests.rs b/storage/overlaybd/src/lsmt/file/tests.rs index 6d24f8bb..c8730f2a 100644 --- a/storage/overlaybd/src/lsmt/file/tests.rs +++ b/storage/overlaybd/src/lsmt/file/tests.rs @@ -535,14 +535,16 @@ async fn rejects_index_count_overflow_before_allocation() { } #[tokio::test] -async fn rejects_index_larger_than_allocation_limit() { +async fn rejects_index_exceeding_decoded_memory_budget() { let temp_dir = TempDir::new().unwrap(); let data = create_sealed_layer(&temp_dir, "oversized-index", 8192, &[]).await; - let count = MAX_INDEX_BYTES as u64 / size_of::() as u64 + 1; + let bytes_per_mapping = + size_of::() as u64 + size_of::() as u64; + let count = MAX_INDEX_MEMORY_BYTES as u64 / bytes_per_mapping + 1; let err = load_index_and_reset_tags(&(data as Arc), HEADER_SIZE, count) .await - .expect_err("oversized index should be rejected before allocation"); - assert_err_contains(&err, "exceeds limit"); + .expect_err("decoded index exceeding the memory budget should be rejected"); + assert_err_contains(&err, "exceeds per-layer limit"); } async fn open_sparse_lsmt_env( diff --git a/storage/overlaybd/src/lsmt/file/types.rs b/storage/overlaybd/src/lsmt/file/types.rs index 968e75b2..3f72ed82 100644 --- a/storage/overlaybd/src/lsmt/file/types.rs +++ b/storage/overlaybd/src/lsmt/file/types.rs @@ -12,7 +12,8 @@ pub(super) const ALIGNMENT: u64 = 512; pub(super) const ALIGNMENT_USIZE: usize = 512; pub(super) const HEADER_SIZE: u64 = 4096; pub(super) const MAX_IO_SIZE: usize = 4 * 1024 * 1024; // 4MB -pub(super) const MAX_INDEX_BYTES: usize = 1024 * 1024 * 1024; // 1 GiB +pub(super) const MAX_INDEX_MEMORY_BYTES: usize = 256 * 1024 * 1024; +pub(super) const MAX_STACK_INDEX_MEMORY_BYTES: usize = 1024 * 1024 * 1024; pub(super) const ALIGNMENT_4K: usize = 4096; pub(super) const INVALID_SEGMENT_OFFSET: u64 = (1 << 50) - 1; pub(super) const COMPACT_ZERO_DETECTION_ENABLED: bool = false; From d02f2f8e27e754f17870b9bab19b1293957f8d9c Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Fri, 31 Jul 2026 23:50:48 +0800 Subject: [PATCH 6/9] fix(overlaybd): bound compact index construction --- storage/overlaybd/src/lsmt/file/helper.rs | 25 ++++++++++++++++++-- storage/overlaybd/src/lsmt/file/readwrite.rs | 1 + storage/overlaybd/src/lsmt/file/tests.rs | 9 +++++-- 3 files changed, 31 insertions(+), 4 deletions(-) diff --git a/storage/overlaybd/src/lsmt/file/helper.rs b/storage/overlaybd/src/lsmt/file/helper.rs index 9cf9e62f..ae35f9bc 100644 --- a/storage/overlaybd/src/lsmt/file/helper.rs +++ b/storage/overlaybd/src/lsmt/file/helper.rs @@ -868,6 +868,20 @@ pub(super) fn validate_index_memory(count: u64) -> Result<()> { Ok(()) } +pub(super) fn reserve_compact_index( + index: &mut Vec, + additional: usize, +) -> Result<()> { + let new_len = index + .len() + .checked_add(additional) + .context("compact index mapping count overflow")?; + validate_index_memory(new_len as u64)?; + index + .try_reserve_exact(additional) + .context("failed to reserve compact index mappings") +} + async fn load_readonly_layer_metadata(file: Arc) -> Result { let file_size = file.size().await?; let header = verify_ht(&file, false, file_size).await?; @@ -1275,9 +1289,11 @@ pub async fn compact_to( let writer = commit_args.writer; let concurrency = commit_args.concurrency.max(1); - writer.write_all_at(&header_buf, 0).await?; + validate_index_memory(mappings.len() as u64)?; + let mut compact_index: Vec = Vec::new(); + reserve_compact_index(&mut compact_index, mappings.len())?; - let mut compact_index: Vec = Vec::with_capacity(mappings.len()); + writer.write_all_at(&header_buf, 0).await?; let mut dest_moffset = HEADER_SIZE / ALIGNMENT; if COMPACT_ZERO_DETECTION_ENABLED { @@ -1287,6 +1303,7 @@ pub async fn compact_to( if m.zeroed { let mut zero = *m; zero.moffset = dest_moffset; + reserve_compact_index(&mut compact_index, 1)?; compact_index.push(zero); continue; } @@ -1297,6 +1314,7 @@ pub async fn compact_to( dest_moffset, ) .await?; + reserve_compact_index(&mut compact_index, entries.len())?; compact_index.extend(entries); dest_moffset += written_blocks; } @@ -1322,6 +1340,7 @@ pub async fn compact_to( if m.zeroed { let mut zero = *m; zero.moffset = dest_moffset; + reserve_compact_index(&mut compact_index, 1)?; compact_index.push(zero); // NOTE: no need to advance dest_moffset, as this is a zero segement, // does not occupy space in dset file. @@ -1375,6 +1394,7 @@ pub async fn compact_to( if concurrency == 1 { for chunk in &chunks { let entries = compact_copy_chunk(src_layers, writer.as_ref(), chunk).await?; + reserve_compact_index(&mut compact_index, entries.len())?; compact_index.extend(entries); } } else { @@ -1398,6 +1418,7 @@ pub async fn compact_to( .collect::>>()?; results.sort_by_key(|(order, _)| *order); for (_, entries) in results { + reserve_compact_index(&mut compact_index, entries.len())?; compact_index.extend(entries); } } diff --git a/storage/overlaybd/src/lsmt/file/readwrite.rs b/storage/overlaybd/src/lsmt/file/readwrite.rs index e9d620dd..e4fb3b6a 100644 --- a/storage/overlaybd/src/lsmt/file/readwrite.rs +++ b/storage/overlaybd/src/lsmt/file/readwrite.rs @@ -880,6 +880,7 @@ impl LSMTFile { if m.tag as usize == self.rw_tag { let mut cm = m; cm.tag = 0; + reserve_compact_index(&mut compact_index, 1)?; compact_index.push(cm); } } diff --git a/storage/overlaybd/src/lsmt/file/tests.rs b/storage/overlaybd/src/lsmt/file/tests.rs index c8730f2a..d2f7a579 100644 --- a/storage/overlaybd/src/lsmt/file/tests.rs +++ b/storage/overlaybd/src/lsmt/file/tests.rs @@ -524,9 +524,14 @@ async fn rejects_index_count_overflow_before_allocation() { let temp_dir = TempDir::new().unwrap(); let data = create_sealed_layer(&temp_dir, "overflow-index", 8192, &[]).await; let count = u64::MAX / size_of::() as u64 + 1; - let err = load_index_and_reset_tags(&(data as Arc), HEADER_SIZE, count) + update_trailer(&data, |trailer| { + trailer.index_size = U64::new(count); + }) + .await; + let err = open_file_ro(data as Arc) .await - .expect_err("overflowing index count should be rejected"); + .err() + .expect("overflowing trailer index count should be rejected"); let message = err.to_string(); assert!( message.contains("overflow") || message.contains("does not fit usize"), From 0359b0ef83f495777ef01c81cc674db8c9f95965 Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Sat, 1 Aug 2026 01:53:01 +0800 Subject: [PATCH 7/9] fix(overlaybd): bound index allocation growth --- storage/overlaybd/src/lsmt/file/helper.rs | 55 ++++++++++++++++---- storage/overlaybd/src/lsmt/file/readwrite.rs | 28 +++++----- storage/overlaybd/src/lsmt/file/stack.rs | 2 +- storage/overlaybd/src/lsmt/file/tests.rs | 28 ++++++++++ 4 files changed, 87 insertions(+), 26 deletions(-) diff --git a/storage/overlaybd/src/lsmt/file/helper.rs b/storage/overlaybd/src/lsmt/file/helper.rs index ae35f9bc..fdc48de0 100644 --- a/storage/overlaybd/src/lsmt/file/helper.rs +++ b/storage/overlaybd/src/lsmt/file/helper.rs @@ -859,6 +859,17 @@ pub(super) fn index_memory_bytes(count: u64) -> Result { .context("index memory size overflow") } +pub(super) fn stack_index_memory_bytes(count: u64) -> Result { + let count = usize::try_from(count).context("index mapping count does not fit usize")?; + // During a stack merge the decoded layer vectors remain live while mappings + // are copied into a BTreeSet and then into the final output vector. Budget + // two SegmentMapping-sized units per BTree entry for payload and node + // overhead, in addition to the decoded and output copies. + count + .checked_mul(size_of::() + 4 * size_of::()) + .context("stack index memory size overflow") +} + pub(super) fn validate_index_memory(count: u64) -> Result<()> { let memory_bytes = index_memory_bytes(count)?; ensure!( @@ -877,11 +888,43 @@ pub(super) fn reserve_compact_index( .checked_add(additional) .context("compact index mapping count overflow")?; validate_index_memory(new_len as u64)?; + if new_len <= index.capacity() { + return Ok(()); + } + + let max_capacity = + MAX_INDEX_MEMORY_BYTES / (size_of::() + size_of::()); + let doubled = index.capacity().saturating_mul(2).max(1); + let target_capacity = new_len.max(doubled.min(max_capacity)); index - .try_reserve_exact(additional) + .try_reserve_exact(target_capacity - index.len()) .context("failed to reserve compact index mappings") } +pub(super) fn serialize_index_with_padding( + compact_index: &[SegmentMapping], + padding_count: usize, +) -> Result> { + let index_entries = compact_index + .len() + .checked_add(padding_count) + .context("padded index mapping count overflow")?; + validate_index_memory(index_entries as u64)?; + let index_bytes_len = index_entries + .checked_mul(size_of::()) + .context("padded index byte length overflow")?; + let mut index_bytes = Vec::new(); + index_bytes + .try_reserve_exact(index_bytes_len) + .context("failed to reserve serialized index buffer")?; + for mapping in compact_index { + let disk_mapping = DiskSegmentMapping::from_memory(mapping); + index_bytes.extend_from_slice(disk_mapping.as_bytes()); + } + append_invalid_index_padding(&mut index_bytes, padding_count); + Ok(index_bytes) +} + async fn load_readonly_layer_metadata(file: Arc) -> Result { let file_size = file.size().await?; let header = verify_ht(&file, false, file_size).await?; @@ -1436,15 +1479,7 @@ pub async fn compact_to( 0 }; let index_size = compact_index.len() as u64 + padding_count as u64; - validate_index_memory(index_size)?; - - let mut index_bytes = Vec::with_capacity(compact_index.len() * size_of::()); - for m in &compact_index { - let dm = DiskSegmentMapping::from_memory(m); - index_bytes.extend_from_slice(dm.as_bytes()); - } - - append_invalid_index_padding(&mut index_bytes, padding_count); + let index_bytes = serialize_index_with_padding(&compact_index, padding_count)?; writer.write_all_at(&index_bytes, index_offset).await?; diff --git a/storage/overlaybd/src/lsmt/file/readwrite.rs b/storage/overlaybd/src/lsmt/file/readwrite.rs index e4fb3b6a..112474bb 100644 --- a/storage/overlaybd/src/lsmt/file/readwrite.rs +++ b/storage/overlaybd/src/lsmt/file/readwrite.rs @@ -875,12 +875,16 @@ impl LSMTFile { idx.lookup(query, &mut mappings); } + let compact_len = mappings + .iter() + .filter(|mapping| mapping.tag as usize == self.rw_tag) + .count(); let mut compact_index: Vec = Vec::new(); + reserve_compact_index(&mut compact_index, compact_len)?; for m in mappings { if m.tag as usize == self.rw_tag { let mut cm = m; cm.tag = 0; - reserve_compact_index(&mut compact_index, 1)?; compact_index.push(cm); } } @@ -890,20 +894,14 @@ impl LSMTFile { // HEADER_SIZE + virtual_size; punching the last virtual block ends // exactly at this boundary, so index bytes cannot overlap virtual data. let index_offset = data_file_size; - validate_index_memory(compact_index.len() as u64)?; - - let mut index_bytes = - Vec::with_capacity(compact_index.len() * size_of::()); - for m in &compact_index { - let dm = DiskSegmentMapping::from_memory(m); - index_bytes.extend_from_slice(dm.as_bytes()); - } - - let remainder = index_bytes.len() % ALIGNMENT_USIZE; - if remainder != 0 { - let pad = ALIGNMENT_USIZE - remainder; - index_bytes.extend(vec![0xff; pad]); - } + let mappings_per_block = ALIGNMENT_USIZE / size_of::(); + let remainder = compact_index.len() % mappings_per_block; + let padding_count = if remainder == 0 { + 0 + } else { + mappings_per_block - remainder + }; + let index_bytes = serialize_index_with_padding(&compact_index, padding_count)?; self.rw_data_file .write_at(index_offset, &index_bytes) diff --git a/storage/overlaybd/src/lsmt/file/stack.rs b/storage/overlaybd/src/lsmt/file/stack.rs index e7a0a2f1..d940fbc9 100644 --- a/storage/overlaybd/src/lsmt/file/stack.rs +++ b/storage/overlaybd/src/lsmt/file/stack.rs @@ -63,7 +63,7 @@ async fn merge_readonly_indexes( ); let total_index_memory = metadata.iter().try_fold(0usize, |total, layer| { total - .checked_add(index_memory_bytes(layer.index_size)?) + .checked_add(stack_index_memory_bytes(layer.index_size)?) .context("stack index memory size overflow") })?; ensure!( diff --git a/storage/overlaybd/src/lsmt/file/tests.rs b/storage/overlaybd/src/lsmt/file/tests.rs index d2f7a579..dcdd2e5e 100644 --- a/storage/overlaybd/src/lsmt/file/tests.rs +++ b/storage/overlaybd/src/lsmt/file/tests.rs @@ -552,6 +552,34 @@ async fn rejects_index_exceeding_decoded_memory_budget() { assert_err_contains(&err, "exceeds per-layer limit"); } +#[test] +fn compact_index_reservation_grows_geometrically_within_budget() { + let mut index = Vec::new(); + reserve_compact_index(&mut index, 1).unwrap(); + index.push(SegmentMapping::default()); + reserve_compact_index(&mut index, 1).unwrap(); + index.push(SegmentMapping::default()); + reserve_compact_index(&mut index, 1).unwrap(); + + assert!(index.capacity() >= 3); + assert!(index.capacity() > index.len()); + validate_index_memory(index.capacity() as u64).unwrap(); +} + +#[test] +fn serialized_index_reserves_space_for_padding() { + let mappings = [SegmentMapping::new(0, 1, 8, false, 0)]; + let bytes = serialize_index_with_padding(&mappings, 2).unwrap(); + + assert_eq!(bytes.len(), 3 * size_of::()); +} + +#[test] +fn stack_index_budget_includes_merge_working_set() { + let count = 1024; + assert!(stack_index_memory_bytes(count).unwrap() > index_memory_bytes(count).unwrap()); +} + async fn open_sparse_lsmt_env( data_file: Arc, lower_layers: Vec>, From 0d1b50fc2687df30125812f6d93cbcd803a63b59 Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Sat, 1 Aug 2026 02:43:27 +0800 Subject: [PATCH 8/9] fix(overlaybd): bound compaction working memory --- storage/overlaybd/src/lsmt/file/helper.rs | 151 +++++++++++++++++----- storage/overlaybd/src/lsmt/file/tests.rs | 9 ++ 2 files changed, 128 insertions(+), 32 deletions(-) diff --git a/storage/overlaybd/src/lsmt/file/helper.rs b/storage/overlaybd/src/lsmt/file/helper.rs index fdc48de0..e006cf9f 100644 --- a/storage/overlaybd/src/lsmt/file/helper.rs +++ b/storage/overlaybd/src/lsmt/file/helper.rs @@ -1065,7 +1065,7 @@ fn push_compact_segment( zero_detected: bool, segment: &mut SegmentMapping, index: &mut Vec, -) { +) -> Result<()> { if zero_detected { segment.zeroed = true; } else { @@ -1075,6 +1075,9 @@ fn push_compact_segment( } *prev_end_blocks += segment.length() as usize; + index + .try_reserve_exact(1) + .context("failed to reserve zero-detection compact mapping")?; index.push(*segment); let next_moffset = segment.mend(); @@ -1083,6 +1086,7 @@ fn push_compact_segment( segment.segment.offset = next_offset; segment.segment.length = 0; segment.moffset = next_moffset; + Ok(()) } // --------------------------------------------------------------------------- @@ -1114,8 +1118,41 @@ struct CompactChunk { /// Total data bytes in this chunk (sum of entry lengths). /// Always `<= writer.buffer_size()`. total_len: usize, - /// Insertion order so parallel results can be merged correctly. - order: usize, +} + +pub(super) fn validate_compaction_memory( + input_mappings: usize, + output_mappings: usize, + active_buffers: usize, + buffer_size: usize, +) -> Result<()> { + let input_bytes = input_mappings + .checked_mul(size_of::()) + .context("compaction input memory size overflow")?; + // Conservatively assume one chunk per output mapping. Charge twice the + // chunk and entry storage because `try_reserve` may grow those vectors + // geometrically, then include a pending result mapping, the final compact + // mapping, and its serialized form even though some lifetimes do not + // overlap. + let output_bytes_per_mapping = 2 * size_of::() + + 2 * size_of::() + + 2 * size_of::() + + size_of::(); + let output_bytes = output_mappings + .checked_mul(output_bytes_per_mapping) + .context("compaction output memory size overflow")?; + let buffer_bytes = active_buffers + .checked_mul(buffer_size) + .context("compaction buffer memory size overflow")?; + let peak_bytes = input_bytes + .checked_add(output_bytes) + .and_then(|bytes| bytes.checked_add(buffer_bytes)) + .context("compaction peak memory size overflow")?; + ensure!( + peak_bytes <= MAX_INDEX_MEMORY_BYTES, + "compaction peak index memory {peak_bytes} exceeds limit {MAX_INDEX_MEMORY_BYTES}" + ); + Ok(()) } /// Read all entries in `chunk` into a single writer-provided buffer and @@ -1130,7 +1167,10 @@ async fn compact_copy_chunk( let mut buf = writer.alloc_buffer().await?; let buf_slice = buf.as_mut().as_mut(); let mut buf_offset = 0usize; - let mut index = Vec::with_capacity(chunk.entries.len()); + let mut index = Vec::new(); + index + .try_reserve_exact(chunk.entries.len()) + .context("failed to reserve compact chunk result mappings")?; for entry in &chunk.entries { let layer_idx = entry.tag as usize; @@ -1200,7 +1240,9 @@ async fn compact_copy_mapping_with_zero_detection( let mut index = Vec::new(); let mut segment = SegmentMapping::new(mapping.offset(), 0, moffset, false, mapping.tag); - let mut data = Vec::with_capacity(buf_size); + let mut data = Vec::new(); + data.try_reserve_exact(buf_size) + .context("failed to reserve zero-detection compact buffer")?; while remaining > 0 { let step = min(remaining as usize, buf_size); @@ -1231,7 +1273,7 @@ async fn compact_copy_mapping_with_zero_detection( false, &mut segment, &mut index, - ); + )?; } segment.segment.length += 1; zero_detected = true; @@ -1247,7 +1289,7 @@ async fn compact_copy_mapping_with_zero_detection( true, &mut segment, &mut index, - ); + )?; } segment.segment.length += 1; zero_detected = false; @@ -1262,7 +1304,7 @@ async fn compact_copy_mapping_with_zero_detection( have_detection && zero_detected, &mut segment, &mut index, - ); + )?; } if !data.is_empty() { @@ -1331,8 +1373,16 @@ pub async fn compact_to( let writer = commit_args.writer; let concurrency = commit_args.concurrency.max(1); + let buf_size = writer.buffer_size(); + ensure!( + buf_size.is_multiple_of(ALIGNMENT_USIZE), + "CompactWriter buffer size {buf_size} not aligned" + ); validate_index_memory(mappings.len() as u64)?; + // `compact_index` reserves this capacity immediately below, so include it + // in the working set before allocating. + validate_compaction_memory(mappings.len(), mappings.len(), 0, buf_size)?; let mut compact_index: Vec = Vec::new(); reserve_compact_index(&mut compact_index, mappings.len())?; @@ -1344,12 +1394,35 @@ pub async fn compact_to( // per mapping, so we cannot pre-chunk or pre-compute offsets. for m in mappings { if m.zeroed { + let output_mappings = compact_index + .len() + .checked_add(1) + .context("compaction output mapping count overflow")?; + validate_compaction_memory( + mappings.len(), + output_mappings.max(mappings.len()), + 0, + buf_size, + )?; let mut zero = *m; zero.moffset = dest_moffset; reserve_compact_index(&mut compact_index, 1)?; compact_index.push(zero); continue; } + // Zero detection can emit at most one run per logical block. + let max_entries = usize::try_from(m.length()) + .context("zero-detection mapping length does not fit usize")?; + let max_output_mappings = compact_index + .len() + .checked_add(max_entries) + .context("compaction output mapping count overflow")?; + validate_compaction_memory( + mappings.len(), + max_output_mappings.max(mappings.len()), + 1, + buf_size, + )?; let (written_blocks, entries) = compact_copy_mapping_with_zero_detection( src_layers, writer.as_ref(), @@ -1370,17 +1443,21 @@ pub async fn compact_to( // // Zeroed mappings don't consume destination space — they are recorded // directly in `compact_index` and skipped during chunking. - let buf_size = writer.buffer_size(); - ensure!( - buf_size.is_multiple_of(ALIGNMENT_USIZE), - "CompactWriter buffer size {buf_size} not aligned" - ); let mut chunks: Vec = Vec::new(); let mut current_chunk: Option = None; - let mut chunk_order = 0usize; + let mut planned_output_mappings = 0usize; for m in mappings { if m.zeroed { + planned_output_mappings = planned_output_mappings + .checked_add(1) + .context("compaction output mapping count overflow")?; + validate_compaction_memory( + mappings.len(), + planned_output_mappings.max(mappings.len()), + concurrency.min(chunks.len() + usize::from(current_chunk.is_some())), + buf_size, + )?; let mut zero = *m; zero.moffset = dest_moffset; reserve_compact_index(&mut compact_index, 1)?; @@ -1395,20 +1472,28 @@ pub async fn compact_to( let mut logical_off = m.offset(); while remaining > 0 { - let chunk = current_chunk.get_or_insert_with(|| { - let c = CompactChunk { - entries: Vec::new(), - dest_moffset, - total_len: 0, - order: chunk_order, - }; - chunk_order += 1; - c + let chunk = current_chunk.get_or_insert_with(|| CompactChunk { + entries: Vec::new(), + dest_moffset, + total_len: 0, }); let space = buf_size - chunk.total_len; let take = remaining.min(space); + planned_output_mappings = planned_output_mappings + .checked_add(1) + .context("compaction output mapping count overflow")?; + validate_compaction_memory( + mappings.len(), + planned_output_mappings.max(mappings.len()), + concurrency.min(chunks.len() + 1), + buf_size, + )?; + chunk + .entries + .try_reserve(1) + .context("failed to reserve compact chunk entry")?; chunk.entries.push(CompactChunkEntry { tag: m.tag, src_offset: src_off, @@ -1424,12 +1509,18 @@ pub async fn compact_to( dest_moffset += blocks_taken; if chunk.total_len >= buf_size { + chunks + .try_reserve(1) + .context("failed to reserve compact chunk")?; chunks.push(current_chunk.take().unwrap()); } } } // Flush remaining partial chunk. if let Some(chunk) = current_chunk.take() { + chunks + .try_reserve(1) + .context("failed to reserve compact chunk")?; chunks.push(chunk); } @@ -1443,7 +1534,7 @@ pub async fn compact_to( } else { let semaphore = Arc::new(tokio::sync::Semaphore::new(concurrency)); let src_layers: Arc<[Arc]> = Arc::from(src_layers.to_vec()); - let mut results: Vec<(usize, Vec)> = stream::iter(chunks) + let mut results = stream::iter(chunks) .map(|chunk| { let sem = semaphore.clone(); let writer = writer.clone(); @@ -1451,16 +1542,12 @@ pub async fn compact_to( async move { let _permit = sem.acquire().await.context("semaphore closed")?; let entries = compact_copy_chunk(&layers, writer.as_ref(), &chunk).await?; - Ok::<_, anyhow::Error>((chunk.order, entries)) + Ok::<_, anyhow::Error>(entries) } }) - .buffer_unordered(concurrency) - .collect::>() - .await - .into_iter() - .collect::>>()?; - results.sort_by_key(|(order, _)| *order); - for (_, entries) in results { + .buffered(concurrency); + while let Some(entries) = results.next().await { + let entries = entries?; reserve_compact_index(&mut compact_index, entries.len())?; compact_index.extend(entries); } diff --git a/storage/overlaybd/src/lsmt/file/tests.rs b/storage/overlaybd/src/lsmt/file/tests.rs index dcdd2e5e..486c674f 100644 --- a/storage/overlaybd/src/lsmt/file/tests.rs +++ b/storage/overlaybd/src/lsmt/file/tests.rs @@ -580,6 +580,15 @@ fn stack_index_budget_includes_merge_working_set() { assert!(stack_index_memory_bytes(count).unwrap() > index_memory_bytes(count).unwrap()); } +#[test] +fn compaction_budget_includes_coexisting_working_sets() { + validate_compaction_memory(1, 1, 1, ALIGNMENT_USIZE).unwrap(); + + let err = validate_compaction_memory(0, usize::MAX, 1, ALIGNMENT_USIZE) + .expect_err("unbounded compaction output should be rejected"); + assert_err_contains(&err, "compaction output memory size overflow"); +} + async fn open_sparse_lsmt_env( data_file: Arc, lower_layers: Vec>, From ebffd4b481cd7078ce31bb52217ad98a96844cfd Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Sat, 1 Aug 2026 03:12:54 +0800 Subject: [PATCH 9/9] fix(overlaybd): reject zero compaction buffers --- storage/overlaybd/src/lsmt/file/helper.rs | 8 ++++---- storage/overlaybd/src/lsmt/file/tests.rs | 13 ++++++++----- 2 files changed, 12 insertions(+), 9 deletions(-) diff --git a/storage/overlaybd/src/lsmt/file/helper.rs b/storage/overlaybd/src/lsmt/file/helper.rs index e006cf9f..dfe82118 100644 --- a/storage/overlaybd/src/lsmt/file/helper.rs +++ b/storage/overlaybd/src/lsmt/file/helper.rs @@ -1234,8 +1234,8 @@ async fn compact_copy_mapping_with_zero_detection( let mut written_bytes = 0u64; let buf_size = writer.buffer_size(); ensure!( - buf_size.is_multiple_of(ALIGNMENT_USIZE), - "CompactWriter buffer size {buf_size} not aligned" + buf_size > 0 && buf_size.is_multiple_of(ALIGNMENT_USIZE), + "CompactWriter buffer size {buf_size} must be non-zero and aligned" ); let mut index = Vec::new(); @@ -1375,8 +1375,8 @@ pub async fn compact_to( let concurrency = commit_args.concurrency.max(1); let buf_size = writer.buffer_size(); ensure!( - buf_size.is_multiple_of(ALIGNMENT_USIZE), - "CompactWriter buffer size {buf_size} not aligned" + buf_size > 0 && buf_size.is_multiple_of(ALIGNMENT_USIZE), + "CompactWriter buffer size {buf_size} must be non-zero and aligned" ); validate_index_memory(mappings.len() as u64)?; diff --git a/storage/overlaybd/src/lsmt/file/tests.rs b/storage/overlaybd/src/lsmt/file/tests.rs index 486c674f..5a5fe0ec 100644 --- a/storage/overlaybd/src/lsmt/file/tests.rs +++ b/storage/overlaybd/src/lsmt/file/tests.rs @@ -556,14 +556,17 @@ async fn rejects_index_exceeding_decoded_memory_budget() { fn compact_index_reservation_grows_geometrically_within_budget() { let mut index = Vec::new(); reserve_compact_index(&mut index, 1).unwrap(); - index.push(SegmentMapping::default()); - reserve_compact_index(&mut index, 1).unwrap(); - index.push(SegmentMapping::default()); + let previous_capacity = index.capacity(); + index.resize(previous_capacity, SegmentMapping::default()); reserve_compact_index(&mut index, 1).unwrap(); - assert!(index.capacity() >= 3); - assert!(index.capacity() > index.len()); + assert!(index.capacity() >= previous_capacity * 2); validate_index_memory(index.capacity() as u64).unwrap(); + + let max_capacity = + MAX_INDEX_MEMORY_BYTES / (size_of::() + size_of::()); + reserve_compact_index(&mut Vec::new(), max_capacity + 1) + .expect_err("reservation beyond the index budget should be rejected"); } #[test]