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 308d5b06 fix(blob): validate entry integrity on reads (#582)
308d5b06 is described below
commit 308d5b06439f1e86d97467f00b783da39285c0fd
Author: QuakeWang <[email protected]>
AuthorDate: Wed Jul 22 11:18:38 2026 +0800
fix(blob): validate entry integrity on reads (#582)
---
crates/paimon/src/arrow/format/blob.rs | 418 ++++++++++++++++++++++++++++-----
1 file changed, 359 insertions(+), 59 deletions(-)
diff --git a/crates/paimon/src/arrow/format/blob.rs
b/crates/paimon/src/arrow/format/blob.rs
index 5b455735..51ea8a97 100644
--- a/crates/paimon/src/arrow/format/blob.rs
+++ b/crates/paimon/src/arrow/format/blob.rs
@@ -108,6 +108,8 @@ pub(crate) enum BlobReadValue {
const BLOB_FOOTER_SIZE: u64 = 5;
const BLOB_FORMAT_VERSION: u8 = 1;
+const BLOB_MAGIC_NUMBER: i32 = 1481511375;
+const BLOB_MAGIC_NUMBER_BYTES: [u8; 4] = BLOB_MAGIC_NUMBER.to_le_bytes();
const BLOB_INLINE_HEADER_SIZE: u64 = 4;
const BLOB_TRAILER_SIZE: u64 = 12;
const BLOB_ENTRY_OVERHEAD: u64 = BLOB_INLINE_HEADER_SIZE + BLOB_TRAILER_SIZE;
@@ -329,8 +331,7 @@ fn plan_blob_reads(
})?;
Ok(match entry {
- BlobEntry::Value(range) if range.start == range.end =>
PlannedBlobRead::Empty,
- BlobEntry::Value(range) =>
PlannedBlobRead::Read(range.clone()),
+ BlobEntry::Value(range) =>
PlannedBlobRead::Entry(blob_entry_range(range)),
BlobEntry::Null => PlannedBlobRead::Null,
BlobEntry::Placeholder => PlannedBlobRead::Placeholder,
})
@@ -346,8 +347,9 @@ async fn fetch_blob_values(
match planned_read {
PlannedBlobRead::Null => Ok(BlobReadValue::Null),
PlannedBlobRead::Placeholder => Ok(BlobReadValue::Placeholder),
- PlannedBlobRead::Empty => Ok(BlobReadValue::Value(Bytes::new())),
- PlannedBlobRead::Read(range) =>
reader.read(range).await.map(BlobReadValue::Value),
+ PlannedBlobRead::Entry(range) => read_blob_entry(reader, range)
+ .await
+ .map(BlobReadValue::Value),
}
}))
.buffered(BLOB_READ_CONCURRENCY)
@@ -355,6 +357,66 @@ async fn fetch_blob_values(
.await
}
+fn blob_entry_range(payload_range: &Range<u64>) -> Range<u64> {
+ payload_range.start - BLOB_INLINE_HEADER_SIZE..payload_range.end +
BLOB_TRAILER_SIZE
+}
+
+async fn read_blob_entry(reader: &dyn FileRead, entry_range: Range<u64>) ->
crate::Result<Bytes> {
+ let expected_entry_length = entry_range.end - entry_range.start;
+ let entry = reader.read(entry_range.clone()).await?;
+ if entry.len() as u64 != expected_entry_length {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Short read for Blob entry range {entry_range:?}: expected
{expected_entry_length} bytes, got {}",
+ entry.len()
+ ),
+ source: None,
+ });
+ }
+
+ let actual_magic = i32::from_le_bytes(
+ entry[..BLOB_INLINE_HEADER_SIZE as usize]
+ .try_into()
+ .unwrap(),
+ );
+ if actual_magic != BLOB_MAGIC_NUMBER {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Invalid Blob entry magic at offset {}: expected
{BLOB_MAGIC_NUMBER}, got {actual_magic}",
+ entry_range.start
+ ),
+ source: None,
+ });
+ }
+
+ let length_offset = entry.len() - BLOB_TRAILER_SIZE as usize;
+ let crc_offset = entry.len() - std::mem::size_of::<u32>();
+ let embedded_length =
i64::from_le_bytes(entry[length_offset..crc_offset].try_into().unwrap());
+ if u64::try_from(embedded_length).ok() != Some(expected_entry_length) {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Blob entry length mismatch at offset {}: index declares
{expected_entry_length}, entry stores {embedded_length}",
+ entry_range.start
+ ),
+ source: None,
+ });
+ }
+
+ let expected_crc =
u32::from_le_bytes(entry[crc_offset..].try_into().unwrap());
+ let actual_crc = crc32fast::hash(&entry[..crc_offset]);
+ if actual_crc != expected_crc {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Blob entry CRC32 mismatch at offset {}: expected
{expected_crc:#010x}, got {actual_crc:#010x}",
+ entry_range.start
+ ),
+ source: None,
+ });
+ }
+
+ Ok(entry.slice(BLOB_INLINE_HEADER_SIZE as usize..length_offset))
+}
+
fn plan_blob_array_reads(
blob_index: &BlobFileIndex,
positions: &[usize],
@@ -392,11 +454,11 @@ async fn fetch_blob_array_values(
PlannedBlobArrayRead::Null => Ok(BlobReadValue::Null),
PlannedBlobArrayRead::Placeholder =>
Ok(BlobReadValue::Placeholder),
PlannedBlobArrayRead::Read(payload_range) => {
- let metadata = read_blob_array_metadata(reader,
payload_range).await?;
if descriptor_mode {
+ let metadata = read_blob_array_metadata(reader,
payload_range).await?;
build_blob_array_descriptors(metadata, file_path)
} else {
- read_inline_blob_array(reader, metadata).await
+ read_inline_blob_array_entry(reader, payload_range).await
}
}
}
@@ -410,6 +472,41 @@ async fn read_blob_array_metadata(
reader: &dyn FileRead,
payload_range: Range<u64>,
) -> crate::Result<BlobArrayMetadata> {
+ let layout = read_blob_array_layout(reader, payload_range).await?;
+ let index_bytes = if layout.element_index_range.is_empty() {
+ Bytes::new()
+ } else {
+ read_blob_array_range(reader, layout.element_index_range.clone(),
"element index").await?
+ };
+
+ decode_blob_array_metadata(layout, index_bytes.as_ref())
+}
+
+async fn read_blob_array_layout(
+ reader: &dyn FileRead,
+ payload_range: Range<u64>,
+) -> crate::Result<BlobArrayLayout> {
+ let payload_length = validate_blob_array_payload_range(&payload_range)?;
+
+ let header_end = payload_range.start + BLOB_ARRAY_HEADER_SIZE;
+ let header = read_blob_array_range(reader,
payload_range.start..header_end, "header").await?;
+
+ let index_length_position = payload_range.end -
BLOB_ARRAY_INDEX_LENGTH_SIZE;
+ let index_length_bytes = read_blob_array_range(
+ reader,
+ index_length_position..payload_range.end,
+ "index length",
+ )
+ .await?;
+ parse_blob_array_layout(
+ payload_range,
+ payload_length,
+ header.as_ref(),
+ index_length_bytes.as_ref(),
+ )
+}
+
+fn validate_blob_array_payload_range(payload_range: &Range<u64>) ->
crate::Result<u64> {
let payload_length = payload_range
.end
.checked_sub(payload_range.start)
@@ -425,9 +522,15 @@ async fn read_blob_array_metadata(
source: None,
});
}
+ Ok(payload_length)
+}
- let header_end = payload_range.start + BLOB_ARRAY_HEADER_SIZE;
- let header = read_blob_array_range(reader,
payload_range.start..header_end, "header").await?;
+fn parse_blob_array_layout(
+ payload_range: Range<u64>,
+ payload_length: u64,
+ header: &[u8],
+ index_length_bytes: &[u8],
+) -> crate::Result<BlobArrayLayout> {
let magic = i32::from_le_bytes(header[..4].try_into().unwrap());
if magic != BLOB_ARRAY_MAGIC_NUMBER {
return Err(Error::DataInvalid {
@@ -453,13 +556,6 @@ async fn read_blob_array_metadata(
});
}
- let index_length_position = payload_range.end -
BLOB_ARRAY_INDEX_LENGTH_SIZE;
- let index_length_bytes = read_blob_array_range(
- reader,
- index_length_position..payload_range.end,
- "index length",
- )
- .await?;
let index_length =
i32::from_le_bytes(index_length_bytes[..4].try_into().unwrap());
let maximum_index_length = payload_length - BLOB_ARRAY_MIN_PAYLOAD_SIZE;
if index_length < 0 || index_length as u64 > maximum_index_length {
@@ -476,29 +572,48 @@ async fn read_blob_array_metadata(
});
}
+ let index_length_position = payload_range.end -
BLOB_ARRAY_INDEX_LENGTH_SIZE;
let index_start = index_length_position - index_length;
- let index_bytes = if index_length == 0 {
- Bytes::new()
- } else {
- read_blob_array_range(reader, index_start..index_length_position,
"element index").await?
- };
- let encoded_lengths =
- decode_delta_varints(index_bytes.as_ref()).map_err(|e|
Error::DataInvalid {
- message: format!("Invalid ARRAY<BLOB> element index: {e}"),
- source: Some(Box::new(e)),
- })?;
- if encoded_lengths.len() != element_count as usize {
+ Ok(BlobArrayLayout {
+ element_count: element_count as usize,
+ element_data_range: payload_range.start +
BLOB_ARRAY_HEADER_SIZE..index_start,
+ element_index_range: index_start..index_length_position,
+ })
+}
+
+fn validate_inline_blob_array_data_length(layout: &BlobArrayLayout) ->
crate::Result<()> {
+ let data_length = layout.element_data_range.end -
layout.element_data_range.start;
+ if data_length > i32::MAX as u64 {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "ARRAY<BLOB> inline element data is too large for Arrow
Binary: {data_length} bytes"
+ ),
+ source: None,
+ });
+ }
+ Ok(())
+}
+
+fn decode_blob_array_metadata(
+ layout: BlobArrayLayout,
+ index_bytes: &[u8],
+) -> crate::Result<BlobArrayMetadata> {
+ let encoded_lengths = decode_delta_varints(index_bytes).map_err(|e|
Error::DataInvalid {
+ message: format!("Invalid ARRAY<BLOB> element index: {e}"),
+ source: Some(Box::new(e)),
+ })?;
+ if encoded_lengths.len() != layout.element_count {
return Err(Error::DataInvalid {
message: format!(
- "ARRAY<BLOB> element count {element_count} does not match
index value count {}",
+ "ARRAY<BLOB> element count {} does not match index value count
{}",
+ layout.element_count,
encoded_lengths.len()
),
source: None,
});
}
- let element_data_range = header_end..index_start;
- let mut remaining_data_length = element_data_range.end -
element_data_range.start;
+ let mut remaining_data_length = layout.element_data_range.end -
layout.element_data_range.start;
let mut element_lengths = Vec::with_capacity(encoded_lengths.len());
for encoded_length in encoded_lengths {
if encoded_length == BLOB_ARRAY_NULL_ELEMENT_LENGTH {
@@ -526,7 +641,7 @@ async fn read_blob_array_metadata(
}
Ok(BlobArrayMetadata {
- element_data_range,
+ element_data_range: layout.element_data_range,
element_lengths,
})
}
@@ -556,24 +671,40 @@ async fn read_blob_array_range(
Ok(bytes)
}
-async fn read_inline_blob_array(
+async fn read_inline_blob_array_entry(
reader: &dyn FileRead,
- metadata: BlobArrayMetadata,
+ payload_range: Range<u64>,
) -> crate::Result<BlobReadValue> {
- let data_length = metadata.element_data_range.end -
metadata.element_data_range.start;
- if data_length > i32::MAX as u64 {
+ let preflight_layout = read_blob_array_layout(reader,
payload_range.clone()).await?;
+ validate_inline_blob_array_data_length(&preflight_layout)?;
+
+ let payload = read_blob_entry(reader,
blob_entry_range(&payload_range)).await?;
+ let payload_length = validate_blob_array_payload_range(&payload_range)?;
+ if payload.len() as u64 != payload_length {
return Err(Error::DataInvalid {
message: format!(
- "ARRAY<BLOB> inline element data is too large for Arrow
Binary: {data_length} bytes"
+ "ARRAY<BLOB> payload length mismatch: expected
{payload_length} bytes, got {}",
+ payload.len()
),
source: None,
});
}
- let data = if data_length == 0 {
- Bytes::new()
- } else {
- read_blob_array_range(reader, metadata.element_data_range, "element
data").await?
- };
+
+ let index_length_position = payload.len() - BLOB_ARRAY_INDEX_LENGTH_SIZE
as usize;
+ let layout = parse_blob_array_layout(
+ payload_range.clone(),
+ payload_length,
+ &payload[..BLOB_ARRAY_HEADER_SIZE as usize],
+ &payload[index_length_position..],
+ )?;
+ validate_inline_blob_array_data_length(&layout)?;
+ let index_start = (layout.element_index_range.start - payload_range.start)
as usize;
+ let index_end = (layout.element_index_range.end - payload_range.start) as
usize;
+ let metadata = decode_blob_array_metadata(layout,
&payload[index_start..index_end])?;
+
+ let data_start = (metadata.element_data_range.start - payload_range.start)
as usize;
+ let data_end = (metadata.element_data_range.end - payload_range.start) as
usize;
+ let data = payload.slice(data_start..data_end);
let mut offset = 0usize;
let mut elements = Vec::with_capacity(metadata.element_lengths.len());
@@ -632,8 +763,7 @@ fn build_blob_array_descriptors(
enum PlannedBlobRead {
Null,
Placeholder,
- Empty,
- Read(Range<u64>),
+ Entry(Range<u64>),
}
#[derive(Debug, Clone)]
@@ -649,6 +779,13 @@ struct BlobArrayMetadata {
element_lengths: Vec<Option<u64>>,
}
+#[derive(Debug)]
+struct BlobArrayLayout {
+ element_count: usize,
+ element_data_range: Range<u64>,
+ element_index_range: Range<u64>,
+}
+
#[derive(Debug, Clone)]
struct BlobFileIndex {
entries: Vec<BlobEntry>,
@@ -974,8 +1111,6 @@ fn decode_varint(bytes: &[u8]) -> crate::Result<(i64,
usize)> {
// --- Blob Format Writer ---
-const BLOB_MAGIC_NUMBER_BYTES: [u8; 4] = 1481511375_i32.to_le_bytes();
-
pub(crate) struct BlobFormatWriter {
writer: Box<dyn FileWrite>,
file_io: Option<crate::io::FileIO>,
@@ -1232,6 +1367,7 @@ mod tests {
use arrow_array::Array;
use bytes::Bytes;
use futures::TryStreamExt;
+ use std::mem::size_of;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Mutex;
use std::time::Duration;
@@ -1366,6 +1502,50 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn test_inline_blob_array_reader_rejects_payload_crc_mismatch() {
+ let payload = build_blob_array_payload(b"helloworld", &[5, -1, 5]);
+ let mut file_bytes =
blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]);
+ let first_element_offset = (BLOB_INLINE_HEADER_SIZE +
BLOB_ARRAY_HEADER_SIZE) as usize;
+ file_bytes[first_element_offset] ^= 0xff;
+
+ let stream = BlobFormatReader::new(String::new(), false)
+ .read_batch_stream(
+ Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))),
+ file_bytes.len() as u64,
+ &blob_array_read_fields(),
+ None,
+ None,
+ None,
+ )
+ .await
+ .unwrap();
+ let error = stream.try_collect::<Vec<_>>().await.unwrap_err();
+
+ assert_data_invalid(error, "CRC32 mismatch");
+ }
+
+ #[tokio::test]
+ async fn
test_inline_blob_array_reader_rejects_oversized_data_before_entry_read() {
+ let element_data_length = i32::MAX as u64 + 1;
+ let element_index = encode_delta_varints_write(&[element_data_length
as i64]);
+ let payload_length =
+ BLOB_ARRAY_MIN_PAYLOAD_SIZE + element_data_length +
element_index.len() as u64;
+ let payload_range = BLOB_INLINE_HEADER_SIZE..BLOB_INLINE_HEADER_SIZE +
payload_length;
+ let reader =
+ BlobArrayPreflightFileRead::new(payload_range.clone(), 1,
element_index.len() as i32);
+
+ let error = read_inline_blob_array_entry(&reader,
payload_range.clone())
+ .await
+ .unwrap_err();
+
+ assert!(
+ !reader.ranges().contains(&blob_entry_range(&payload_range)),
+ "oversized inline ARRAY<BLOB> must be rejected before reading the
complete entry"
+ );
+ assert_data_invalid(error, "too large");
+ }
+
#[tokio::test]
async fn
test_blob_array_reader_builds_exact_descriptors_without_payload_reads() {
let file_path = "file:///tmp/blob-array.blob";
@@ -1433,21 +1613,17 @@ mod tests {
.await
.unwrap();
- assert_eq!(
- reader
- .ranges()
- .iter()
- .filter(|range| **range == (13..23))
- .count(),
- 1
- );
- assert!(
- !reader
- .ranges()
- .iter()
- .any(|range| (13..23).contains(&range.start) && *range !=
(13..23)),
- "inline mode must not issue per-element payload reads"
- );
+ let element_data_range = 13..23;
+ let expected_entry_end =
+ BLOB_ENTRY_OVERHEAD + build_blob_array_payload(b"helloworld", &[5,
-1, 5]).len() as u64;
+ let overlapping_reads = reader
+ .ranges()
+ .into_iter()
+ .filter(|range| {
+ range.start < element_data_range.end &&
element_data_range.start < range.end
+ })
+ .collect::<Vec<_>>();
+ assert_eq!(overlapping_reads, vec![0..expected_entry_end]);
}
#[tokio::test]
@@ -1876,6 +2052,52 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn test_blob_reader_rejects_invalid_entry_magic() {
+ let payload = b"hello";
+ let mut file_bytes =
blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]);
+ file_bytes[..BLOB_INLINE_HEADER_SIZE as
usize].copy_from_slice(&0_i32.to_le_bytes());
+ rewrite_first_blob_entry_crc(&mut file_bytes, payload.len());
+
+ let error = read_scalar_blob_values(file_bytes).await.unwrap_err();
+ assert_data_invalid(error, "magic");
+ }
+
+ #[tokio::test]
+ async fn test_blob_reader_rejects_mismatched_entry_length() {
+ let payload = b"hello";
+ let mut file_bytes =
blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]);
+ let length_offset = BLOB_INLINE_HEADER_SIZE as usize + payload.len();
+ let mismatched_length = payload.len() as i64 + BLOB_ENTRY_OVERHEAD as
i64 + 1;
+ file_bytes[length_offset..length_offset + size_of::<i64>()]
+ .copy_from_slice(&mismatched_length.to_le_bytes());
+ rewrite_first_blob_entry_crc(&mut file_bytes, payload.len());
+
+ let error = read_scalar_blob_values(file_bytes).await.unwrap_err();
+ assert_data_invalid(error, "length mismatch");
+ }
+
+ #[tokio::test]
+ async fn test_blob_reader_rejects_payload_crc_mismatch() {
+ let payload = b"hello";
+ let mut file_bytes =
blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]);
+ file_bytes[BLOB_INLINE_HEADER_SIZE as usize] ^= 0xff;
+
+ let error = read_scalar_blob_values(file_bytes).await.unwrap_err();
+ assert_data_invalid(error, "CRC32 mismatch");
+ }
+
+ #[tokio::test]
+ async fn test_blob_reader_rejects_empty_entry_crc_mismatch() {
+ let payload = b"";
+ let mut file_bytes =
blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]);
+ let crc_offset = BLOB_INLINE_HEADER_SIZE as usize + size_of::<i64>();
+ file_bytes[crc_offset] ^= 0xff;
+
+ let error = read_scalar_blob_values(file_bytes).await.unwrap_err();
+ assert_data_invalid(error, "CRC32 mismatch");
+ }
+
#[test]
fn test_varint_encode_decode_roundtrip() {
let values = vec![21, -1, 0, i64::MAX, i64::MIN + 1, 127, -128, 300,
-300];
@@ -1963,6 +2185,34 @@ mod tests {
.unwrap_err()
}
+ async fn read_scalar_blob_values(file_bytes: Vec<u8>) ->
crate::Result<Vec<Option<Vec<u8>>>> {
+ let batches = BlobFormatReader::new(String::new(), false)
+ .read_batch_stream(
+ Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))),
+ file_bytes.len() as u64,
+ &[DataField::new(
+ 0,
+ "payload".to_string(),
+ DataType::Blob(BlobType::new()),
+ )],
+ None,
+ None,
+ None,
+ )
+ .await?
+ .try_collect::<Vec<_>>()
+ .await?;
+ Ok(batches.iter().flat_map(collect_binary_values).collect())
+ }
+
+ fn rewrite_first_blob_entry_crc(file_bytes: &mut [u8], payload_length:
usize) {
+ let crc_offset = BLOB_INLINE_HEADER_SIZE as usize + payload_length +
size_of::<i64>();
+ let mut hasher = crc32fast::Hasher::new();
+ hasher.update(&file_bytes[..crc_offset]);
+ file_bytes[crc_offset..crc_offset + size_of::<u32>()]
+ .copy_from_slice(&hasher.finalize().to_le_bytes());
+ }
+
fn assert_data_invalid(error: Error, expected_message: &str) {
assert!(
matches!(error, Error::DataInvalid { message, .. } if
message.contains(expected_message)),
@@ -2050,4 +2300,54 @@ mod tests {
Ok(self.bytes.slice(range.start as usize..range.end as usize))
}
}
+
+ struct BlobArrayPreflightFileRead {
+ payload_range: Range<u64>,
+ element_count: i32,
+ index_length: i32,
+ ranges: Mutex<Vec<Range<u64>>>,
+ }
+
+ impl BlobArrayPreflightFileRead {
+ fn new(payload_range: Range<u64>, element_count: i32, index_length:
i32) -> Self {
+ Self {
+ payload_range,
+ element_count,
+ index_length,
+ ranges: Mutex::new(Vec::new()),
+ }
+ }
+
+ fn ranges(&self) -> Vec<Range<u64>> {
+ self.ranges.lock().unwrap().clone()
+ }
+ }
+
+ #[async_trait::async_trait]
+ impl FileRead for BlobArrayPreflightFileRead {
+ async fn read(&self, range: Range<u64>) -> crate::Result<Bytes> {
+ self.ranges.lock().unwrap().push(range.clone());
+
+ let header_range =
+ self.payload_range.start..self.payload_range.start +
BLOB_ARRAY_HEADER_SIZE;
+ if range == header_range {
+ let mut header = Vec::with_capacity(BLOB_ARRAY_HEADER_SIZE as
usize);
+
header.extend_from_slice(&BLOB_ARRAY_MAGIC_NUMBER.to_le_bytes());
+ header.push(BLOB_ARRAY_VERSION);
+ header.extend_from_slice(&self.element_count.to_le_bytes());
+ return Ok(Bytes::from(header));
+ }
+
+ let index_length_range =
+ self.payload_range.end -
BLOB_ARRAY_INDEX_LENGTH_SIZE..self.payload_range.end;
+ if range == index_length_range {
+ return
Ok(Bytes::copy_from_slice(&self.index_length.to_le_bytes()));
+ }
+
+ Err(Error::UnexpectedError {
+ message: format!("Unexpected sparse Blob test read:
{range:?}"),
+ source: None,
+ })
+ }
+ }
}