This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new 7bcdf1aa fix(table): reject invalid global index metadata (#578)
7bcdf1aa is described below
commit 7bcdf1aa25b600bfaf4c571f59da7a9cace2ef67
Author: XiaoHongbo <[email protected]>
AuthorDate: Wed Jul 22 18:18:35 2026 +0800
fix(table): reject invalid global index metadata (#578)
---
crates/paimon/src/btree/meta.rs | 91 ++++++++++---
crates/paimon/src/table/global_index_scanner.rs | 173 +++++++++++++++++++++---
2 files changed, 229 insertions(+), 35 deletions(-)
diff --git a/crates/paimon/src/btree/meta.rs b/crates/paimon/src/btree/meta.rs
index ea511ea3..c8c0bbe8 100644
--- a/crates/paimon/src/btree/meta.rs
+++ b/crates/paimon/src/btree/meta.rs
@@ -32,6 +32,28 @@ const FORMAT_VERSION_WITH_NULL_FLAGS: u8 = 1;
const FIRST_KEY_IS_NULL: u8 = 1;
const LAST_KEY_IS_NULL: u8 = 1 << 1;
+fn invalid_meta(message: &'static str) -> io::Error {
+ io::Error::new(io::ErrorKind::InvalidData, message)
+}
+
+fn read_key(data: &[u8], pos: &mut usize) -> io::Result<Vec<u8>> {
+ let remaining = data
+ .get(*pos..)
+ .ok_or_else(|| invalid_meta("BTreeIndexMeta key offset out of
bounds"))?;
+ let length_bytes = remaining
+ .get(..4)
+ .ok_or_else(|| invalid_meta("BTreeIndexMeta key length is
truncated"))?;
+ let length = i32::from_le_bytes(length_bytes.try_into().unwrap());
+ let length = usize::try_from(length)
+ .map_err(|_| invalid_meta("BTreeIndexMeta key length is negative"))?;
+ let key = remaining
+ .get(4..)
+ .and_then(|bytes| bytes.get(..length))
+ .ok_or_else(|| invalid_meta("BTreeIndexMeta key data is truncated"))?;
+ *pos += 4 + length;
+ Ok(key.to_vec())
+}
+
/// Index meta for each BTree index file.
#[derive(Debug, Clone)]
pub struct BTreeIndexMeta {
@@ -166,26 +188,22 @@ impl BTreeIndexMeta {
let mut pos = 0;
- let fk_len = i32::from_le_bytes(data[pos..pos +
4].try_into().unwrap()) as usize;
- pos += 4;
- let mut first_key = {
- let key = data[pos..pos + fk_len].to_vec();
- pos += fk_len;
- Some(key)
- };
-
- let lk_len = i32::from_le_bytes(data[pos..pos +
4].try_into().unwrap()) as usize;
- pos += 4;
- let mut last_key = {
- let key = data[pos..pos + lk_len].to_vec();
- pos += lk_len;
- Some(key)
- };
-
- let has_nulls = data[pos] == 1;
+ let mut first_key = Some(read_key(data, &mut pos)?);
+ let mut last_key = Some(read_key(data, &mut pos)?);
+ let has_nulls = *data
+ .get(pos)
+ .ok_or_else(|| invalid_meta("BTreeIndexMeta has_nulls flag is
missing"))?
+ == 1;
pos += 1;
- if data.len().saturating_sub(pos) >= 2 {
+ let trailer_len = data.len().saturating_sub(pos);
+ if trailer_len == 1 {
+ return Err(invalid_meta(
+ "BTreeIndexMeta null flags trailer is truncated",
+ ));
+ }
+
+ if trailer_len >= 2 {
let format_version = data[pos];
pos += 1;
if format_version == FORMAT_VERSION_WITH_NULL_FLAGS {
@@ -197,7 +215,10 @@ impl BTreeIndexMeta {
last_key = None;
}
}
- } else if fk_len == 0 && lk_len == 0 && has_nulls {
+ } else if first_key.as_ref().is_some_and(Vec::is_empty)
+ && last_key.as_ref().is_some_and(Vec::is_empty)
+ && has_nulls
+ {
first_key = None;
last_key = None;
}
@@ -253,4 +274,36 @@ mod tests {
assert!(!decoded.has_nulls);
assert!(!decoded.only_nulls());
}
+
+ #[test]
+ fn test_meta_rejects_invalid_key_lengths() {
+ let mut negative_first = vec![0; 9];
+ negative_first[..4].copy_from_slice(&(-1i32).to_le_bytes());
+ let mut truncated_first = vec![0; 9];
+ truncated_first[..4].copy_from_slice(&10i32.to_le_bytes());
+ let mut negative_last = vec![0; 9];
+ negative_last[4..8].copy_from_slice(&(-1i32).to_le_bytes());
+ let mut truncated_last = vec![0; 9];
+ truncated_last[4..8].copy_from_slice(&10i32.to_le_bytes());
+
+ for encoded in [
+ negative_first,
+ truncated_first,
+ negative_last,
+ truncated_last,
+ ] {
+ let error = BTreeIndexMeta::deserialize(&encoded).unwrap_err();
+ assert_eq!(error.kind(), io::ErrorKind::InvalidData);
+ }
+ }
+
+ #[test]
+ fn test_meta_rejects_truncated_null_flags_trailer() {
+ let meta = BTreeIndexMeta::new(Some(Vec::new()), Some(Vec::new()),
true);
+ let mut encoded = meta.serialize();
+ assert_eq!(encoded.pop(), Some(0));
+
+ let error = BTreeIndexMeta::deserialize(&encoded).unwrap_err();
+ assert_eq!(error.kind(), io::ErrorKind::InvalidData);
+ }
}
diff --git a/crates/paimon/src/table/global_index_scanner.rs
b/crates/paimon/src/table/global_index_scanner.rs
index adfc628a..c535db7f 100644
--- a/crates/paimon/src/table/global_index_scanner.rs
+++ b/crates/paimon/src/table/global_index_scanner.rs
@@ -33,7 +33,7 @@ use crate::spec::{
Predicate, PredicateOperator,
};
use crate::table::{DeletionFile, RowRange, Table};
-use crate::Result;
+use crate::{Error, Result};
use roaring::RoaringTreemap;
use std::cmp::Ordering;
use std::collections::{HashMap, HashSet};
@@ -151,7 +151,7 @@ impl OpenedGlobalIndexReader {
impl GlobalIndexScanner {
/// Create a scanner from index manifest entries.
- /// Returns `None` if there are no global index entries.
+ /// Returns `Ok(None)` if there are no global index entries.
pub(crate) fn create(
file_io: &FileIO,
table_path: &str,
@@ -159,7 +159,7 @@ impl GlobalIndexScanner {
bitmap_fallback_scan_max_size: i64,
index_entries: &[IndexManifestEntry],
schema_fields: &[DataField],
- ) -> Option<Self> {
+ ) -> Result<Option<Self>> {
let mut entries_by_field: std::collections::HashMap<i32,
Vec<GlobalIndexEntry>> =
std::collections::HashMap::new();
let mut coverage_by_field: HashMap<i32, Vec<RowRange>> =
HashMap::new();
@@ -172,16 +172,37 @@ impl GlobalIndexScanner {
else {
continue;
};
- let global_meta = match &entry.index_file.global_index_meta {
- Some(m) => m,
- None => continue,
- };
+ let global_meta =
+ entry
+ .index_file
+ .global_index_meta
+ .as_ref()
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Missing global index metadata for sorted index
file '{}'",
+ entry.index_file.file_name
+ ),
+ source: None,
+ })?;
- let sorted_meta = global_meta
+ let index_meta = global_meta
.index_meta
.as_ref()
- .and_then(|bytes| BTreeIndexMeta::deserialize(bytes).ok())
- .unwrap_or_else(|| BTreeIndexMeta::new(None, None, false));
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Missing sorted global index metadata for file '{}'",
+ entry.index_file.file_name
+ ),
+ source: None,
+ })?;
+ let sorted_meta =
+ BTreeIndexMeta::deserialize(index_meta).map_err(|error|
Error::DataInvalid {
+ message: format!(
+ "Invalid sorted global index metadata for file '{}'",
+ entry.index_file.file_name
+ ),
+ source: Some(Box::new(error)),
+ })?;
let resolved = GlobalIndexEntry {
file_name: entry.index_file.file_name.clone(),
@@ -216,10 +237,10 @@ impl GlobalIndexScanner {
}
if entries_by_field.is_empty() {
- return None;
+ return Ok(None);
}
- Some(Self {
+ Ok(Some(Self {
file_io: file_io.clone(),
table_path: table_path.trim_end_matches('/').to_string(),
btree_fallback_scan_max_size,
@@ -228,7 +249,7 @@ impl GlobalIndexScanner {
coverage_by_field,
schema_fields: schema_fields.to_vec(),
reader_cache: Mutex::new(HashMap::new()),
- })
+ }))
}
/// Evaluate a predicate against the global indexes and return matching
row ranges.
@@ -487,7 +508,10 @@ impl GlobalIndexScanner {
let mut reader = if (between_matches &&
between_evaluated_for_entry)
|| !matching_predicates.is_empty()
{
- Some(self.open_reader_for_entry(entry, data_type).await?)
+ Some(
+ self.open_reader_for_entry(entry, &entry.meta, data_type)
+ .await?,
+ )
} else {
None
};
@@ -589,11 +613,12 @@ impl GlobalIndexScanner {
async fn open_reader_for_entry(
&self,
entry: &GlobalIndexEntry,
+ meta: &BTreeIndexMeta,
data_type: &DataType,
) -> Result<OpenedGlobalIndexReader> {
match entry.index_type {
GlobalIndexFileKind::BTree => {
- self.get_or_open_reader(&entry.file_name, &entry.meta,
data_type)
+ self.get_or_open_reader(&entry.file_name, meta, data_type)
.await
}
GlobalIndexFileKind::Bitmap => self
@@ -1199,7 +1224,7 @@ pub(crate) async fn evaluate_global_index(
evaluation.bitmap_fallback_scan_max_size,
evaluation.index_entries,
evaluation.schema_fields,
- ) {
+ )? {
Some(s) => s,
None => return Ok(None),
};
@@ -1544,6 +1569,7 @@ mod tests {
&entries,
&fields,
)
+ .expect("create scanner")
.expect("scanner");
let ranges = scanner
@@ -1571,6 +1597,7 @@ mod tests {
&entries,
&fields,
)
+ .expect("create scanner")
.expect("scanner");
let ranges = scanner
@@ -1598,6 +1625,7 @@ mod tests {
&entries,
&fields,
)
+ .expect("create scanner")
.expect("scanner");
let ranges = scanner
@@ -1632,6 +1660,7 @@ mod tests {
&entries,
&fields,
)
+ .expect("create scanner")
.expect("scanner");
let predicate = Predicate::and(vec![int_eq("id", 0, 7),
int_eq("value", 1, 8)]);
@@ -1655,6 +1684,7 @@ mod tests {
&entries,
&fields,
)
+ .expect("create scanner")
.expect("scanner");
let predicate = Predicate::and(vec![int_eq("id", 0, 7),
int_eq("value", 1, 8)]);
@@ -1684,6 +1714,7 @@ mod tests {
&[entry],
&fields,
)
+ .expect("create scanner")
.expect("scanner");
let ranges = scanner
@@ -1747,6 +1778,116 @@ mod tests {
assert_eq!(ranges, vec![RowRange::new(25, 25)]);
}
+ #[tokio::test]
+ async fn test_missing_index_meta_returns_error() {
+ let (file_io, table_path, file_name, tmp) =
+ setup_testdata_table("btree_int_100_no_compress.bin");
+ let second_file_name = "btree_int_100_no_compress_2.bin";
+ std::fs::copy(
+ tmp.path().join("index").join(&file_name),
+ tmp.path().join("index").join(second_file_name),
+ )
+ .unwrap();
+ let meta = BTreeIndexMeta::new(Some(le_int_key(0)),
Some(le_int_key(198)), false);
+ let valid_entry = make_global_index_entry(&file_name, 1, 0, 99, &meta);
+ let mut invalid_entry = make_global_index_entry(second_file_name, 1,
100, 199, &meta);
+ invalid_entry
+ .index_file
+ .global_index_meta
+ .as_mut()
+ .unwrap()
+ .index_meta = None;
+
+ let error = evaluate_global_index_fast(
+ &file_io,
+ &table_path,
+ &[valid_entry, invalid_entry],
+ &[int_eq("id", 0, 50)],
+ &int_schema_fields(),
+ )
+ .await
+ .expect_err("missing sorted index metadata must fail the scan");
+
+ assert!(matches!(
+ error,
+ crate::Error::DataInvalid { message, .. }
+ if message.contains(second_file_name)
+ ));
+ }
+
+ #[tokio::test]
+ async fn test_invalid_index_meta_returns_error() {
+ let (file_io, table_path, file_name, tmp) =
+ setup_testdata_table("btree_int_100_no_compress.bin");
+ let second_file_name = "btree_int_100_no_compress_2.bin";
+ std::fs::copy(
+ tmp.path().join("index").join(&file_name),
+ tmp.path().join("index").join(second_file_name),
+ )
+ .unwrap();
+ let meta = BTreeIndexMeta::new(Some(le_int_key(0)),
Some(le_int_key(198)), false);
+ let valid_entry = make_global_index_entry(&file_name, 1, 0, 99, &meta);
+ let mut invalid_entry = make_global_index_entry(second_file_name, 1,
100, 199, &meta);
+ let mut invalid_meta = vec![0; 9];
+ invalid_meta[..4].copy_from_slice(&10i32.to_le_bytes());
+ invalid_entry
+ .index_file
+ .global_index_meta
+ .as_mut()
+ .unwrap()
+ .index_meta = Some(invalid_meta);
+
+ let error = evaluate_global_index_fast(
+ &file_io,
+ &table_path,
+ &[valid_entry, invalid_entry],
+ &[int_eq("id", 0, 50)],
+ &int_schema_fields(),
+ )
+ .await
+ .expect_err("invalid sorted index metadata must fail the scan");
+
+ assert!(matches!(
+ error,
+ crate::Error::DataInvalid {
+ message,
+ source: Some(_),
+ } if message.contains(second_file_name)
+ ));
+ }
+
+ #[tokio::test]
+ async fn test_missing_global_index_meta_returns_error() {
+ let (file_io, table_path, file_name, tmp) =
+ setup_testdata_table("btree_int_100_no_compress.bin");
+ let second_file_name = "btree_int_100_no_compress_2.bin";
+ std::fs::copy(
+ tmp.path().join("index").join(&file_name),
+ tmp.path().join("index").join(second_file_name),
+ )
+ .unwrap();
+ let meta = BTreeIndexMeta::new(Some(le_int_key(0)),
Some(le_int_key(198)), false);
+ let valid_entry = make_global_index_entry(&file_name, 1, 0, 99, &meta);
+ let mut invalid_entry = make_global_index_entry(second_file_name, 1,
100, 199, &meta);
+ invalid_entry.index_file.global_index_meta = None;
+
+ let error = evaluate_global_index_fast(
+ &file_io,
+ &table_path,
+ &[valid_entry, invalid_entry],
+ &[int_eq("id", 0, 50)],
+ &int_schema_fields(),
+ )
+ .await
+ .expect_err("missing global index metadata must fail the scan");
+
+ assert!(matches!(
+ error,
+ crate::Error::DataInvalid { message, .. }
+ if message.contains(second_file_name)
+ ));
+ }
+
#[tokio::test]
async fn test_evaluate_java_bitmap_golden_index_eq_and_null() {
let data_type =
DataType::VarChar(crate::spec::VarCharType::string_type());