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

spetz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git


The following commit(s) were added to refs/heads/master by this push:
     new cb57ce2bd fix(connectors): validate protobuf field lengths 
consistently (#4069)
cb57ce2bd is described below

commit cb57ce2bd569d415bdd02b4638ba76c134472138
Author: Piotr Gankiewicz <[email protected]>
AuthorDate: Mon Sep 7 07:56:10 2026 +0200

    fix(connectors): validate protobuf field lengths consistently (#4069)
    
    Schema-based protobuf decoding did not reliably reject invalid
    field lengths. Use shared checked conversion, addition, and bounds
    validation for known, preserved unknown, and skipped fields.
    
    Add regression coverage for invalid lengths and valid decoding
    with an in-memory schema.
---
 core/connectors/sdk/src/decoders/proto.rs | 135 ++++++++++++++++++++++++++----
 1 file changed, 120 insertions(+), 15 deletions(-)

diff --git a/core/connectors/sdk/src/decoders/proto.rs 
b/core/connectors/sdk/src/decoders/proto.rs
index 1b44bbd84..9e97ba52b 100644
--- a/core/connectors/sdk/src/decoders/proto.rs
+++ b/core/connectors/sdk/src/decoders/proto.rs
@@ -398,11 +398,7 @@ impl ProtoStreamDecoder {
             }
             2 => {
                 let (length, mut new_cursor) = self.parse_simple_varint(data, 
cursor)?;
-                let end_cursor = new_cursor + length as usize;
-
-                if end_cursor > data.len() {
-                    return Err(Error::InvalidProtobufPayload);
-                }
+                let end_cursor = Self::length_delimited_end(new_cursor, 
length, data.len())?;
 
                 let field_data = &data[new_cursor..end_cursor];
                 new_cursor = end_cursor;
@@ -451,11 +447,7 @@ impl ProtoStreamDecoder {
             }
             2 => {
                 let (length, mut new_cursor) = self.parse_simple_varint(data, 
cursor)?;
-                let end_cursor = new_cursor + length as usize;
-
-                if end_cursor > data.len() {
-                    return Err(Error::InvalidProtobufPayload);
-                }
+                let end_cursor = Self::length_delimited_end(new_cursor, 
length, data.len())?;
 
                 let field_data = &data[new_cursor..end_cursor];
                 new_cursor = end_cursor;
@@ -475,11 +467,7 @@ impl ProtoStreamDecoder {
             }
             2 => {
                 let (length, new_cursor) = self.parse_simple_varint(data, 
cursor)?;
-                let end_cursor = new_cursor + length as usize;
-
-                if end_cursor > data.len() {
-                    return Err(Error::InvalidProtobufPayload);
-                }
+                let end_cursor = Self::length_delimited_end(new_cursor, 
length, data.len())?;
 
                 Ok(end_cursor)
             }
@@ -487,6 +475,14 @@ impl ProtoStreamDecoder {
         }
     }
 
+    fn length_delimited_end(cursor: usize, length: u64, data_len: usize) -> 
Result<usize, Error> {
+        usize::try_from(length)
+            .ok()
+            .and_then(|length| cursor.checked_add(length))
+            .filter(|end_cursor| *end_cursor <= data_len)
+            .ok_or(Error::InvalidProtobufPayload)
+    }
+
     fn apply_field_transformations(&self, payload: Payload) -> Result<Payload, 
Error> {
         if let Some(mappings) = &self.config.field_mappings {
             match payload {
@@ -828,4 +824,113 @@ mod tests {
             "Should fallback gracefully when schema loading fails"
         );
     }
+
+    fn one_field_schema_decoder(preserve_unknown_fields: bool) -> 
ProtoStreamDecoder {
+        let file_descriptor_set = prost_types::FileDescriptorSet {
+            file: vec![prost_types::FileDescriptorProto {
+                name: Some("one_field.proto".to_string()),
+                package: Some("test".to_string()),
+                message_type: vec![prost_types::DescriptorProto {
+                    name: Some("OneField".to_string()),
+                    field: vec![prost_types::FieldDescriptorProto {
+                        name: Some("name".to_string()),
+                        number: Some(1),
+                        r#type: Some(Type::String as i32),
+                        ..Default::default()
+                    }],
+                    ..Default::default()
+                }],
+                ..Default::default()
+            }],
+        };
+        let decoder = ProtoStreamDecoder::new(ProtoConfig {
+            descriptor_set: Some(file_descriptor_set.encode_to_vec()),
+            message_type: Some("test.OneField".to_string()),
+            use_any_wrapper: false,
+            preserve_unknown_fields,
+            ..ProtoConfig::default()
+        });
+        assert!(
+            decoder.message_descriptor.is_some(),
+            "schema must be loaded for these tests to exercise the field 
parser"
+        );
+        decoder
+    }
+
+    fn decode_with_loaded_schema(
+        decoder: &ProtoStreamDecoder,
+        payload: &[u8],
+    ) -> Result<Payload, Error> {
+        decoder.decode_with_message_descriptor(
+            payload,
+            decoder.message_descriptor.as_ref().unwrap(),
+            decoder.file_descriptor_set.as_ref().unwrap(),
+        )
+    }
+
+    #[test]
+    fn decode_should_succeed_given_valid_message_with_loaded_schema() {
+        let decoder = one_field_schema_decoder(false);
+        let payload = vec![0x0a, 0x03, b'a', b'b', b'c'];
+
+        let result = decoder.decode(payload);
+
+        let Ok(Payload::Json(simd_json::OwnedValue::Object(map))) = result 
else {
+            panic!("Expected JSON object");
+        };
+        assert_eq!(
+            map.get("name"),
+            Some(&simd_json::OwnedValue::String("abc".to_string()))
+        );
+    }
+
+    #[test]
+    fn decode_should_fail_given_known_field_length_that_overflows_usize() {
+        let decoder = one_field_schema_decoder(false);
+        let payload = vec![
+            0x0a, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0x01,
+        ];
+
+        let result = decode_with_loaded_schema(&decoder, &payload);
+
+        assert!(matches!(result, Err(Error::InvalidProtobufPayload)));
+        assert!(decoder.decode(payload).is_err());
+    }
+
+    #[test]
+    fn 
decode_should_fail_given_unknown_preserved_field_length_that_overflows_usize() {
+        let decoder = one_field_schema_decoder(true);
+        let payload = vec![
+            0x9a, 0x06, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 
0x01,
+        ];
+
+        let result = decode_with_loaded_schema(&decoder, &payload);
+
+        assert!(matches!(result, Err(Error::InvalidProtobufPayload)));
+        assert!(decoder.decode(payload).is_err());
+    }
+
+    #[test]
+    fn decode_should_fail_given_skipped_field_length_that_wraps_cursor() {
+        let decoder = one_field_schema_decoder(false);
+        let payload = vec![
+            0x9a, 0x06, 0xf4, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 
0x01,
+        ];
+
+        let result = decode_with_loaded_schema(&decoder, &payload);
+
+        assert!(matches!(result, Err(Error::InvalidProtobufPayload)));
+        assert!(decoder.decode(payload).is_err());
+    }
+
+    #[test]
+    fn decode_should_fail_given_field_length_past_end_of_payload() {
+        let decoder = one_field_schema_decoder(false);
+        let payload = vec![0x0a, 0x10, b'a'];
+
+        let result = decode_with_loaded_schema(&decoder, &payload);
+
+        assert!(matches!(result, Err(Error::InvalidProtobufPayload)));
+        assert!(decoder.decode(payload).is_err());
+    }
 }

Reply via email to