Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
91 changes: 72 additions & 19 deletions crates/paimon/src/btree/meta.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand All @@ -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;
}
Expand Down Expand Up @@ -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);
}
}
105 changes: 91 additions & 14 deletions crates/paimon/src/table/global_index_scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ struct GlobalIndexEntry {
index_type: GlobalIndexFileKind,
file_size: i64,
row_range_start: i64,
meta: BTreeIndexMeta,
meta: Option<BTreeIndexMeta>,
}

#[derive(Clone, Copy, Debug, PartialEq, Eq)]
Expand Down Expand Up @@ -172,16 +172,12 @@ 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()?;

let sorted_meta = global_meta
.index_meta
.as_ref()
.and_then(|bytes| BTreeIndexMeta::deserialize(bytes).ok())
.unwrap_or_else(|| BTreeIndexMeta::new(None, None, false));
.and_then(|bytes| BTreeIndexMeta::deserialize(bytes).ok());

let resolved = GlobalIndexEntry {
file_name: entry.index_file.file_name.clone(),
Expand Down Expand Up @@ -364,6 +360,14 @@ impl GlobalIndexScanner {
entries: &[GlobalIndexEntry],
predicates: &[(PredicateOperator, &[Datum], &DataType)],
) -> Result<Option<Vec<RowRange>>> {
let Some(metas) = entries
.iter()
.map(|entry| entry.meta.as_ref())
.collect::<Option<Vec<_>>>()
else {
return Ok(None);
};

// Try to detect between pattern and split into (between, remaining)
let (between, remaining) = extract_between(predicates);

Expand Down Expand Up @@ -391,9 +395,9 @@ impl GlobalIndexScanner {
let predicate_matches: Vec<Vec<bool>> = pruning_info
.iter()
.map(|(op, cmp, serialized)| {
entries
metas
.iter()
.map(|entry| entry.meta.may_match(*op, serialized, cmp))
.map(|meta| meta.may_match(*op, serialized, cmp))
.collect()
})
.collect();
Expand All @@ -411,9 +415,9 @@ impl GlobalIndexScanner {
let cmp = make_key_comparator(b.data_type);
let from_key = serialize_datum(b.from, b.data_type);
let to_key = serialize_datum(b.to, b.data_type);
entries
metas
.iter()
.map(|entry| entry.meta.may_match_between(&from_key, &to_key, &cmp))
.map(|meta| meta.may_match_between(&from_key, &to_key, &cmp))
.collect()
}
None => Vec::new(),
Expand All @@ -422,7 +426,7 @@ impl GlobalIndexScanner {
.as_ref()
.map(|_| self.fallback_scan_plan(entries, &between_matches_by_entry));

for (entry_idx, entry) in entries.iter().enumerate() {
for (entry_idx, (entry, meta)) in entries.iter().zip(metas).enumerate() {
// Also check if between range may match
let between_matches = between
.as_ref()
Expand Down Expand Up @@ -487,7 +491,7 @@ 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, meta, data_type).await?)
} else {
None
};
Expand Down Expand Up @@ -589,11 +593,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
Expand Down Expand Up @@ -1747,6 +1752,78 @@ mod tests {
assert_eq!(ranges, vec![RowRange::new(25, 25)]);
}

#[tokio::test]
async fn test_missing_or_invalid_index_meta_falls_back() {
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 mut invalid_meta = vec![0; 9];
invalid_meta[..4].copy_from_slice(&10i32.to_le_bytes());

for index_meta in [None, Some(invalid_meta)] {
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 = index_meta;

let result = evaluate_global_index_fast(
&file_io,
&table_path,
&[valid_entry, invalid_entry],
&[int_eq("id", 0, 50)],
&int_schema_fields(),
)
.await
.unwrap();

assert!(
result.is_none(),
"missing or invalid index metadata must fall back to the normal table scan"
);
}
}

#[tokio::test]
async fn test_missing_global_index_meta_falls_back() {
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 result = evaluate_global_index_fast(
&file_io,
&table_path,
&[valid_entry, invalid_entry],
&[int_eq("id", 0, 50)],
&int_schema_fields(),
)
.await
.unwrap();

assert!(
result.is_none(),
"missing global index metadata must fall back to the normal table scan"
);
}

#[tokio::test]
async fn test_evaluate_java_bitmap_golden_index_eq_and_null() {
let data_type = DataType::VarChar(crate::spec::VarCharType::string_type());
Expand Down
Loading