kafrust_protocol/codec/
mod.rs1mod 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_arrays_before_allocating_above_the_configured_limit() {
107 let limits = DecodeLimits::new().with_max_array_elements(2);
108
109 let array_length = 3_i32.to_be_bytes();
110 let mut decoder = Decoder::with_limits(&array_length, limits);
111 assert_eq!(
112 decoder.read_array("test array", |_| Ok(())).unwrap_err(),
113 Error::LimitExceeded {
114 kind: "test array",
115 actual: 3,
116 max: 2,
117 }
118 );
119
120 let mut decoder = Decoder::with_limits(&[4], limits);
121 assert_eq!(
122 decoder
123 .read_compact_array("compact test array", |_| Ok(()))
124 .unwrap_err(),
125 Error::LimitExceeded {
126 kind: "compact test array",
127 actual: 3,
128 max: 2,
129 }
130 );
131
132 let mut decoder = Decoder::with_limits(&[3], limits);
133 assert_eq!(
134 decoder.read_tagged_fields().unwrap_err(),
135 Error::LimitExceeded {
136 kind: "tagged fields",
137 actual: 3,
138 max: 2,
139 }
140 );
141 }
142}