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> { + 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..7663a640 100644 --- a/crates/paimon/src/table/global_index_scanner.rs +++ b/crates/paimon/src/table/global_index_scanner.rs @@ -82,7 +82,7 @@ struct GlobalIndexEntry { index_type: GlobalIndexFileKind, file_size: i64, row_range_start: i64, - meta: BTreeIndexMeta, + meta: Option, } #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -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(), @@ -364,6 +360,14 @@ impl GlobalIndexScanner { entries: &[GlobalIndexEntry], predicates: &[(PredicateOperator, &[Datum], &DataType)], ) -> Result>> { + let Some(metas) = entries + .iter() + .map(|entry| entry.meta.as_ref()) + .collect::>>() + else { + return Ok(None); + }; + // Try to detect between pattern and split into (between, remaining) let (between, remaining) = extract_between(predicates); @@ -391,9 +395,9 @@ impl GlobalIndexScanner { let predicate_matches: Vec> = 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(); @@ -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(), @@ -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() @@ -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 }; @@ -589,11 +593,12 @@ impl GlobalIndexScanner { async fn open_reader_for_entry( &self, entry: &GlobalIndexEntry, + meta: &BTreeIndexMeta, data_type: &DataType, ) -> Result { 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 @@ -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());