Skip to main content

kafrust_protocol/codec/
mod.rs

1mod decode;
2mod encode;
3
4pub use decode::{DecodeLimits, Decoder, TaggedField};
5pub use encode::Encoder;
6
7#[cfg(test)]
8#[allow(clippy::unwrap_used)]
9mod tests {
10    use super::{DecodeLimits, Decoder, Encoder};
11    use crate::Error;
12
13    #[test]
14    fn encodes_and_decodes_fixed_width_primitives() {
15        let mut encoder = Encoder::new();
16        encoder.write_bool(true);
17        encoder.write_i16(0x1234);
18        encoder.write_i32(0x1234_5678);
19        encoder.write_i64(0x0102_0304_0506_0708);
20        encoder.write_f64(123.5);
21
22        let bytes = encoder.into_bytes();
23        assert_eq!(
24            bytes,
25            [
26                1, 0x12, 0x34, 0x12, 0x34, 0x56, 0x78, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07,
27                0x08, 0x40, 0x5e, 0xe0, 0x00, 0x00, 0x00, 0x00, 0x00,
28            ]
29        );
30
31        let mut decoder = Decoder::new(&bytes);
32        assert!(decoder.read_bool().unwrap());
33        assert_eq!(decoder.read_i16().unwrap(), 0x1234);
34        assert_eq!(decoder.read_i32().unwrap(), 0x1234_5678);
35        assert_eq!(decoder.read_i64().unwrap(), 0x0102_0304_0506_0708);
36        assert_eq!(decoder.read_f64().unwrap(), 123.5);
37        assert!(decoder.is_empty());
38    }
39
40    #[test]
41    fn encodes_and_decodes_nullable_strings_and_bytes() {
42        let mut encoder = Encoder::new();
43        encoder.write_nullable_string(Some("topic")).unwrap();
44        encoder.write_nullable_string(None).unwrap();
45        encoder.write_nullable_bytes(Some(&[1, 2, 3])).unwrap();
46        encoder.write_nullable_bytes(None).unwrap();
47
48        let bytes = encoder.into_bytes();
49        let mut decoder = Decoder::new(&bytes);
50        assert_eq!(
51            decoder.read_nullable_string().unwrap(),
52            Some("topic".to_owned())
53        );
54        assert_eq!(decoder.read_nullable_string().unwrap(), None);
55        assert_eq!(decoder.read_nullable_bytes().unwrap(), Some(vec![1, 2, 3]));
56        assert_eq!(decoder.read_nullable_bytes().unwrap(), None);
57        assert!(decoder.is_empty());
58    }
59
60    #[test]
61    fn encodes_and_decodes_compact_types_and_empty_tags() {
62        let mut encoder = Encoder::new();
63        encoder.write_unsigned_varint(300);
64        encoder.write_compact_string("kafka").unwrap();
65        encoder.write_compact_nullable_string(None).unwrap();
66        encoder.write_compact_bytes(&[9, 8]).unwrap();
67        encoder.write_empty_tagged_fields();
68
69        let bytes = encoder.into_bytes();
70        let mut decoder = Decoder::new(&bytes);
71        assert_eq!(decoder.read_unsigned_varint().unwrap(), 300);
72        assert_eq!(decoder.read_compact_string().unwrap(), "kafka");
73        assert_eq!(decoder.read_compact_nullable_string().unwrap(), None);
74        assert_eq!(decoder.read_compact_bytes().unwrap(), vec![9, 8]);
75        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
76        assert!(decoder.is_empty());
77    }
78
79    #[test]
80    fn decodes_signed_varints() {
81        let mut decoder = Decoder::new(&[0x02, 0x01, 0x04, 0x03]);
82
83        assert_eq!(decoder.read_varint().unwrap(), 1);
84        assert_eq!(decoder.read_varint().unwrap(), -1);
85        assert_eq!(decoder.read_varlong().unwrap(), 2);
86        assert_eq!(decoder.read_varlong().unwrap(), -2);
87        assert!(decoder.is_empty());
88    }
89
90    #[test]
91    fn rejects_invalid_bool_and_short_input() {
92        let mut decoder = Decoder::new(&[2]);
93        assert_eq!(decoder.read_bool(), Err(Error::InvalidBool(2)));
94
95        let mut decoder = Decoder::new(&[0]);
96        assert!(matches!(
97            decoder.read_i16(),
98            Err(Error::UnexpectedEof {
99                needed: 2,
100                remaining: 1
101            })
102        ));
103    }
104
105    #[test]
106    fn rejects_trailing_bytes_when_decoding_finishes() {
107        let mut decoder = Decoder::new(&[0, 0xa5]);
108        assert_eq!(decoder.read_i8().unwrap(), 0);
109        assert_eq!(decoder.finish(), Err(Error::TrailingBytes { remaining: 1 }));
110
111        let mut decoder = Decoder::new(&[0]);
112        assert_eq!(decoder.read_i8().unwrap(), 0);
113        assert_eq!(decoder.finish(), Ok(()));
114    }
115
116    #[test]
117    fn rejects_arrays_before_allocating_above_the_configured_limit() {
118        let limits = DecodeLimits::new().with_max_array_elements(2);
119
120        let array_length = 3_i32.to_be_bytes();
121        let mut decoder = Decoder::with_limits(&array_length, limits);
122        assert_eq!(
123            decoder.read_array("test array", |_| Ok(())).unwrap_err(),
124            Error::LimitExceeded {
125                kind: "test array",
126                actual: 3,
127                max: 2,
128            }
129        );
130
131        let mut decoder = Decoder::with_limits(&[4], limits);
132        assert_eq!(
133            decoder
134                .read_compact_array("compact test array", |_| Ok(()))
135                .unwrap_err(),
136            Error::LimitExceeded {
137                kind: "compact test array",
138                actual: 3,
139                max: 2,
140            }
141        );
142
143        let mut decoder = Decoder::with_limits(&[3], limits);
144        assert_eq!(
145            decoder.read_tagged_fields().unwrap_err(),
146            Error::LimitExceeded {
147                kind: "tagged fields",
148                actual: 3,
149                max: 2,
150            }
151        );
152    }
153}