This is an automated email from the ASF dual-hosted git repository.

martin-g pushed a commit to branch track-allocated-bytes-during-decoding
in repository https://gitbox.apache.org/repos/asf/avro-rs.git

commit fecaf4bd8f0d5341867da8e3166f019fab683168
Author: Martin Tzvetanov Grigorov <[email protected]>
AuthorDate: Tue Aug 25 14:17:39 2026 +0300

    fix: Nested collections multiply the 512MB allocation budget arbitrarily
    
    decode.rs::decode_internal (for Schema::Array & Map): safe_collection_len 
bounds each individual Vec/HashMap against max_allocation_bytes, but no state 
tracks the running total across the decoded value tree. Elements that cost zero 
wire bytes (null, empty records, zero-size fixed) let each ~7-byte inner 
array<null> block materialize ~512MB of live Value slots, so array<array<null>> 
from an attacker OCF header holds K x 512MB simultaneously from ~130 datum 
bytes.
    
    Reported-by: Security scans
---
 avro/src/decode.rs               | 135 ++++++++++++++++++++++++++++++++++++---
 avro/src/reader/block.rs         |   4 +-
 avro/src/reader/datum.rs         |  11 +++-
 avro/src/reader/single_object.rs |   4 +-
 4 files changed, 140 insertions(+), 14 deletions(-)

diff --git a/avro/src/decode.rs b/avro/src/decode.rs
index d78da71..483c5e9 100644
--- a/avro/src/decode.rs
+++ b/avro/src/decode.rs
@@ -24,7 +24,10 @@ use crate::{
     error::Details,
     schema::{DecimalSchema, EnumSchema, FixedSchema, Name, RecordSchema, 
ResolvedSchema, Schema},
     types::Value,
-    util::{safe_collection_len, safe_len, zag_i32, zag_i64},
+    util::{
+        DEFAULT_MAX_ALLOCATION_BYTES, max_allocation_bytes, 
safe_collection_len, safe_len, zag_i32,
+        zag_i64,
+    },
 };
 use std::{
     borrow::Borrow,
@@ -68,10 +71,56 @@ fn decode_seq_len<R: Read>(reader: &mut R) -> 
AvroResult<usize> {
     )
 }
 
+/// Per-datum decoding state.
+///
+/// Tracks the cumulative number of bytes allocated on behalf of a single
+/// datum, so that nested collections cannot multiply the allocation budget:
+/// every allocation performed while decoding one datum is debited from a
+/// shared budget of [`max_allocation_bytes`] bytes, instead of each
+/// collection only being checked in isolation.
+pub(crate) struct DecodeContext {
+    /// Bytes still available for allocations while decoding the current datum.
+    remaining_budget: usize,
+}
+
+impl DecodeContext {
+    /// Create a fresh context. Call once per datum.
+    pub(crate) fn new() -> Self {
+        Self {
+            remaining_budget: 
max_allocation_bytes(DEFAULT_MAX_ALLOCATION_BYTES),
+        }
+    }
+
+    /// Debit `bytes` from the per-datum allocation budget, erroring when the
+    /// cumulative allocations for this datum would exceed it.
+    fn debit(&mut self, bytes: usize) -> AvroResult<()> {
+        match self.remaining_budget.checked_sub(bytes) {
+            Some(remaining) => {
+                self.remaining_budget = remaining;
+                Ok(())
+            }
+            None => Err(Details::MemoryAllocation {
+                desired: Some(bytes),
+                maximum: max_allocation_bytes(DEFAULT_MAX_ALLOCATION_BYTES),
+            }
+            .into()),
+        }
+    }
+
+    /// Debit the budget for `items` collection elements of type `T`.
+    fn debit_items<T>(&mut self, items: usize) -> AvroResult<()> {
+        let bytes = items
+            .checked_mul(size_of::<T>())
+            .ok_or(Details::IntegerOverflow)?;
+        self.debit(bytes)
+    }
+}
+
 /// Decode a `Value` from avro format given its `Schema`.
 pub fn decode<R: Read>(schema: &Schema, reader: &mut R) -> AvroResult<Value> {
     let rs = ResolvedSchema::try_from(schema)?;
-    decode_internal(schema, rs.get_names(), None, reader)
+    let mut ctx = DecodeContext::new();
+    decode_internal(schema, rs.get_names(), None, reader, &mut ctx)
 }
 
 pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
@@ -79,6 +128,7 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
     names: &HashMap<Name, S>,
     enclosing_namespace: NamespaceRef,
     reader: &mut R,
+    ctx: &mut DecodeContext,
 ) -> AvroResult<Value> {
     match schema {
         Schema::Null => Ok(Value::Null),
@@ -106,27 +156,28 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
                     names,
                     enclosing_namespace,
                     reader,
+                    ctx,
                 )? {
                     Value::Fixed(_, bytes) => 
Ok(Value::Decimal(Decimal::from(bytes))),
                     value => Err(Details::FixedValue(value).into()),
                 }
             }
             InnerDecimalSchema::Bytes => {
-                match decode_internal(&Schema::Bytes, names, 
enclosing_namespace, reader)? {
+                match decode_internal(&Schema::Bytes, names, 
enclosing_namespace, reader, ctx)? {
                     Value::Bytes(bytes) => 
Ok(Value::Decimal(Decimal::from(bytes))),
                     value => Err(Details::BytesValue(value).into()),
                 }
             }
         },
         Schema::BigDecimal => {
-            match decode_internal(&Schema::Bytes, names, enclosing_namespace, 
reader)? {
+            match decode_internal(&Schema::Bytes, names, enclosing_namespace, 
reader, ctx)? {
                 Value::Bytes(bytes) => 
deserialize_big_decimal(&bytes).map(Value::BigDecimal),
                 value => Err(Details::BytesValue(value).into()),
             }
         }
         Schema::Uuid(UuidSchema::String) => {
             let Value::String(string) =
-                decode_internal(&Schema::String, names, enclosing_namespace, 
reader)?
+                decode_internal(&Schema::String, names, enclosing_namespace, 
reader, ctx)?
             else {
                 // decoding a String can also return a Null, indicating EOF
                 return Err(Error::new(Details::ReadBytes(std::io::Error::from(
@@ -138,7 +189,7 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
         }
         Schema::Uuid(UuidSchema::Bytes) => {
             let Value::Bytes(bytes) =
-                decode_internal(&Schema::Bytes, names, enclosing_namespace, 
reader)?
+                decode_internal(&Schema::Bytes, names, enclosing_namespace, 
reader, ctx)?
             else {
                 unreachable!(
                     "decode_internal(Schema::Bytes) can only return a 
Value::Bytes or an error"
@@ -153,6 +204,7 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
                 names,
                 enclosing_namespace,
                 reader,
+                ctx,
             )?
             else {
                 unreachable!(
@@ -205,12 +257,14 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
         }
         Schema::Bytes => {
             let len = decode_len(reader)?;
+            ctx.debit(len)?;
             let mut buf = vec![0u8; len];
             reader.read_exact(&mut buf).map_err(Details::ReadBytes)?;
             Ok(Value::Bytes(buf))
         }
         Schema::String => {
             let len = decode_len(reader)?;
+            ctx.debit(len)?;
             let mut buf = vec![0u8; len];
             match reader.read_exact(&mut buf) {
                 Ok(_) => Ok(Value::String(
@@ -226,6 +280,10 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
             }
         }
         Schema::Fixed(FixedSchema { size, .. }) => {
+            // The size is schema-declared, not wire-declared, but the schema
+            // itself may be attacker-supplied (e.g. an OCF header), so it must
+            // be debited from the allocation budget like any other length.
+            ctx.debit(*size)?;
             let mut buf = vec![0u8; *size];
             reader
                 .read_exact(&mut buf)
@@ -247,6 +305,10 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
                     .checked_add(len)
                     .ok_or(Details::IntegerOverflow)?;
                 safe_collection_len::<Value>(total)?;
+                // Elements can be arbitrarily cheap on the wire (e.g. null),
+                // so also debit the per-datum budget: nested collections must
+                // not multiply the allocation budget.
+                ctx.debit_items::<Value>(len)?;
                 // Use reserve_exact as reserve can allocate more than needed 
defeating the purpose
                 // of the previous check
                 items.reserve_exact(len);
@@ -256,6 +318,7 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
                         names,
                         enclosing_namespace,
                         reader,
+                        ctx,
                     )?);
                 }
             }
@@ -279,13 +342,21 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
                     .checked_add(len)
                     .ok_or(Details::IntegerOverflow)?;
                 safe_collection_len::<(String, Value)>(total)?;
+                // See the Array arm: nested collections share one budget.
+                ctx.debit_items::<(String, Value)>(len)?;
 
                 items.reserve(len);
                 for _ in 0..len {
-                    match decode_internal(&Schema::String, names, 
enclosing_namespace, reader)? {
+                    match decode_internal(&Schema::String, names, 
enclosing_namespace, reader, ctx)?
+                    {
                         Value::String(key) => {
-                            let value =
-                                decode_internal(&inner.types, names, 
enclosing_namespace, reader)?;
+                            let value = decode_internal(
+                                &inner.types,
+                                names,
+                                enclosing_namespace,
+                                reader,
+                                ctx,
+                            )?;
                             items.insert(key, value);
                         }
                         value => return 
Err(Details::MapKeyType(value.into()).into()),
@@ -304,7 +375,7 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
                         index,
                         num_variants: variants.len(),
                     })?;
-                let value = decode_internal(variant, names, 
enclosing_namespace, reader)?;
+                let value = decode_internal(variant, names, 
enclosing_namespace, reader, ctx)?;
                 Ok(Value::Union(index as u32, Box::new(value)))
             }
             Err(Details::ReadVariableIntegerBytes(io_err)) => {
@@ -318,9 +389,13 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
         },
         Schema::Record(RecordSchema { name, fields, .. }) => {
             let fully_qualified_name = 
name.fully_qualified_name(enclosing_namespace);
+            // Records can consume zero wire bytes (e.g. all-null fields), so
+            // debit the budget for the field vector and the cloned names.
+            ctx.debit_items::<(String, Value)>(fields.len())?;
             // Benchmarks indicate ~10% improvement using this method.
             let mut items = Vec::with_capacity(fields.len());
             for field in fields {
+                ctx.debit(field.name.len())?;
                 // TODO: This clone is also expensive. See if we can do away 
with it...
                 items.push((
                     field.name.clone(),
@@ -329,6 +404,7 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
                         names,
                         fully_qualified_name.namespace(),
                         reader,
+                        ctx,
                     )?,
                 ));
             }
@@ -339,6 +415,9 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
                 let index = usize::try_from(raw_index)
                     .map_err(|e| Details::ConvertI32ToUsize(e, raw_index))?;
                 if (0..symbols.len()).contains(&index) {
+                    // Cloning the symbol allocates without consuming wire
+                    // bytes, so it counts against the per-datum budget.
+                    ctx.debit(symbols[index].len())?;
                     let symbol = symbols[index].clone();
                     Value::Enum(raw_index as u32, symbol)
                 } else {
@@ -360,6 +439,7 @@ pub(crate) fn decode_internal<R: Read, S: Borrow<Schema>>(
                     names,
                     fully_qualified_name.namespace(),
                     reader,
+                    ctx,
                 )
             } else {
                 
Err(Details::SchemaResolutionError(fully_qualified_name.into_owned()).into())
@@ -472,6 +552,41 @@ mod tests {
         Ok(())
     }
 
+    #[test]
+    fn avro_rs_639_test_nested_collections_share_one_allocation_budget() -> 
TestResult {
+        use crate::util::{DEFAULT_MAX_ALLOCATION_BYTES, max_allocation_bytes};
+
+        // Each inner array<null> block passes the per-collection check on its
+        // own, but the shared per-datum budget must reject the cumulative
+        // total: elements of type null cost zero wire bytes, so without a
+        // cumulative budget a handful of ~10-byte inner arrays would pin an
+        // unbounded multiple of the allocation limit in memory at once.
+        let budget = max_allocation_bytes(DEFAULT_MAX_ALLOCATION_BYTES);
+        let inner_count = (budget / 2) / size_of::<Value>() + 1;
+
+        let inner_arrays_count = 2;
+        let mut payload = Vec::new();
+        // Outer array: a single block declaring two inner arrays.
+        crate::util::zig_i64(inner_arrays_count, &mut payload)?;
+        for _ in 0..inner_arrays_count {
+            payload.extend(create_block(inner_count as i64));
+        }
+        // Outer array terminator.
+        payload.push(0x00);
+
+        let result = decode(
+            &Schema::array(Schema::array(Schema::Null).build()).build(),
+            &mut payload.as_slice(),
+        );
+
+        assert!(
+            result.is_err(),
+            "nested collections must share one allocation budget, got 
{result:?}"
+        );
+
+        Ok(())
+    }
+
     #[test]
     fn test_decode_array_int64_min_block_count_is_rejected() -> TestResult {
         // i64::MIN as a negative block count cannot be negated (checked_neg
diff --git a/avro/src/reader/block.rs b/avro/src/reader/block.rs
index 5e9ae1b..8c9041a 100644
--- a/avro/src/reader/block.rs
+++ b/avro/src/reader/block.rs
@@ -27,7 +27,7 @@ use serde_json::from_slice;
 
 use crate::{
     AvroResult, Codec, Error,
-    decode::{decode, decode_internal},
+    decode::{DecodeContext, decode, decode_internal},
     error::Details,
     schema::{Names, Schema, resolve_names, resolve_names_with_schemata},
     serde::deser_schema::{Config, SchemaAwareDeserializer},
@@ -195,11 +195,13 @@ impl<'r, R: Read> Block<'r, R> {
         let mut block_bytes = &self.buf[self.buf_idx..];
         let b_original = block_bytes.len();
 
+        let mut ctx = DecodeContext::new();
         let item = decode_internal(
             &self.writer_schema,
             &self.names_refs,
             None,
             &mut block_bytes,
+            &mut ctx,
         )?;
         let item = match read_schema {
             Some(schema) => item.resolve(schema)?,
diff --git a/avro/src/reader/datum.rs b/avro/src/reader/datum.rs
index 4f0c31e..36c8fc5 100644
--- a/avro/src/reader/datum.rs
+++ b/avro/src/reader/datum.rs
@@ -22,7 +22,7 @@ use serde::de::DeserializeOwned;
 
 use crate::{
     AvroResult, AvroSchema, Schema,
-    decode::decode_internal,
+    decode::{DecodeContext, decode_internal},
     schema::{ResolvedOwnedSchema, ResolvedSchema},
     serde::deser_schema::{Config, SchemaAwareDeserializer},
     types::Value,
@@ -130,7 +130,14 @@ impl<'s, S: generic_datum_reader_builder::State> 
GenericDatumReaderBuilder<'s, S
 impl<'s> GenericDatumReader<'s> {
     /// Read a Avro datum from the reader.
     pub fn read_value<R: Read>(&self, reader: &mut R) -> AvroResult<Value> {
-        let value = decode_internal(self.writer, self.resolved.get_names(), 
None, reader)?;
+        let mut ctx = DecodeContext::new();
+        let value = decode_internal(
+            self.writer,
+            self.resolved.get_names(),
+            None,
+            reader,
+            &mut ctx,
+        )?;
         if let Some((reader, resolved)) = &self.reader {
             value.resolve_internal(reader, resolved.get_names(), None, None)
         } else {
diff --git a/avro/src/reader/single_object.rs b/avro/src/reader/single_object.rs
index d65ca2f..bf91238 100644
--- a/avro/src/reader/single_object.rs
+++ b/avro/src/reader/single_object.rs
@@ -22,7 +22,7 @@ use serde::de::DeserializeOwned;
 
 use crate::{
     AvroResult, AvroSchema, Schema,
-    decode::decode_internal,
+    decode::{DecodeContext, decode_internal},
     error::Details,
     headers::{HeaderBuilder, RabinFingerprintHeader},
     schema::ResolvedOwnedSchema,
@@ -60,11 +60,13 @@ impl GenericSingleObjectReader {
 impl GenericSingleObjectReader {
     pub fn read_value<R: Read>(&self, reader: &mut R) -> AvroResult<Value> {
         self.read_header(reader)?;
+        let mut ctx = DecodeContext::new();
         decode_internal(
             self.write_schema.get_root_schema(),
             self.write_schema.get_names(),
             None,
             reader,
+            &mut ctx,
         )
     }
 

Reply via email to