Skip to main content

kafrust_protocol/api/
elect_leaders.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 43;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct ElectLeadersRequestV0 {
9    pub correlation_id: i32,
10    pub client_id: Option<String>,
11    pub topics: Option<Vec<ElectLeadersTopicV0>>,
12    pub timeout_ms: i32,
13}
14
15impl ElectLeadersRequestV0 {
16    pub fn encode(&self) -> Result<Vec<u8>> {
17        let mut encoder = Encoder::new();
18        RequestHeader {
19            api_key: API_KEY,
20            api_version: 0,
21            correlation_id: self.correlation_id,
22            client_id: self.client_id.clone(),
23        }
24        .encode_v1(&mut encoder)?;
25        encode_legacy_topics(&mut encoder, self.topics.as_deref())?;
26        encoder.write_i32(self.timeout_ms);
27        Ok(encoder.into_bytes())
28    }
29}
30
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct ElectLeadersRequestV1 {
33    pub correlation_id: i32,
34    pub client_id: Option<String>,
35    pub election_type: i8,
36    pub topics: Option<Vec<ElectLeadersTopicV0>>,
37    pub timeout_ms: i32,
38}
39
40impl ElectLeadersRequestV1 {
41    pub fn encode(&self) -> Result<Vec<u8>> {
42        let mut encoder = Encoder::new();
43        RequestHeader {
44            api_key: API_KEY,
45            api_version: 1,
46            correlation_id: self.correlation_id,
47            client_id: self.client_id.clone(),
48        }
49        .encode_v1(&mut encoder)?;
50        encoder.write_i8(self.election_type);
51        encode_legacy_topics(&mut encoder, self.topics.as_deref())?;
52        encoder.write_i32(self.timeout_ms);
53        Ok(encoder.into_bytes())
54    }
55}
56
57#[derive(Debug, Clone, PartialEq, Eq)]
58pub struct ElectLeadersRequestV2 {
59    pub correlation_id: i32,
60    pub client_id: Option<String>,
61    pub election_type: i8,
62    pub topics: Option<Vec<ElectLeadersTopicV0>>,
63    pub timeout_ms: i32,
64}
65
66impl ElectLeadersRequestV2 {
67    pub fn encode(&self) -> Result<Vec<u8>> {
68        let mut encoder = Encoder::new();
69        RequestHeader {
70            api_key: API_KEY,
71            api_version: 2,
72            correlation_id: self.correlation_id,
73            client_id: self.client_id.clone(),
74        }
75        .encode_v2(&mut encoder)?;
76        encoder.write_i8(self.election_type);
77        encoder.write_compact_array(self.topics.as_deref(), |encoder, topic| {
78            encoder.write_compact_string(&topic.name)?;
79            encoder.write_compact_array(Some(&topic.partitions), |encoder, partition| {
80                encoder.write_i32(*partition);
81                Ok(())
82            })?;
83            encoder.write_empty_tagged_fields();
84            Ok(())
85        })?;
86        encoder.write_i32(self.timeout_ms);
87        encoder.write_empty_tagged_fields();
88        Ok(encoder.into_bytes())
89    }
90}
91
92#[derive(Debug, Clone, PartialEq, Eq)]
93pub struct ElectLeadersTopicV0 {
94    pub name: String,
95    pub partitions: Vec<i32>,
96}
97
98#[derive(Debug, Clone, PartialEq, Eq)]
99pub struct ElectLeadersResponseV0 {
100    pub throttle_time_ms: i32,
101    pub results: Vec<ElectLeadersTopicResultV0>,
102}
103
104impl ElectLeadersResponseV0 {
105    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
106        Ok(Self {
107            throttle_time_ms: decoder.read_i32()?,
108            results: decoder
109                .read_array("elect leaders results", decode_legacy_topic_result)?
110                .unwrap_or_default(),
111        })
112    }
113}
114
115#[derive(Debug, Clone, PartialEq, Eq)]
116pub struct ElectLeadersResponseV1 {
117    pub throttle_time_ms: i32,
118    pub error_code: i16,
119    pub results: Vec<ElectLeadersTopicResultV0>,
120}
121
122impl ElectLeadersResponseV1 {
123    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
124        Ok(Self {
125            throttle_time_ms: decoder.read_i32()?,
126            error_code: decoder.read_i16()?,
127            results: decoder
128                .read_array("elect leaders results", decode_legacy_topic_result)?
129                .unwrap_or_default(),
130        })
131    }
132}
133
134#[derive(Debug, Clone, PartialEq, Eq)]
135pub struct ElectLeadersResponseV2 {
136    pub throttle_time_ms: i32,
137    pub error_code: i16,
138    pub results: Vec<ElectLeadersTopicResultV0>,
139}
140
141impl ElectLeadersResponseV2 {
142    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
143        let throttle_time_ms = decoder.read_i32()?;
144        let error_code = decoder.read_i16()?;
145        let results = decoder
146            .read_compact_array("elect leaders results", |decoder| {
147                let name = decoder.read_compact_string()?;
148                let partitions = decoder
149                    .read_compact_array("elect leaders partition results", |decoder| {
150                        let partition_index = decoder.read_i32()?;
151                        let error_code = decoder.read_i16()?;
152                        let error_message = decoder.read_compact_nullable_string()?;
153                        decoder.read_tagged_fields()?;
154                        Ok(ElectLeadersPartitionResultV0 {
155                            partition_index,
156                            error_code,
157                            error_message,
158                        })
159                    })?
160                    .unwrap_or_default();
161                decoder.read_tagged_fields()?;
162                Ok(ElectLeadersTopicResultV0 { name, partitions })
163            })?
164            .unwrap_or_default();
165        decoder.read_tagged_fields()?;
166        Ok(Self {
167            throttle_time_ms,
168            error_code,
169            results,
170        })
171    }
172}
173
174#[derive(Debug, Clone, PartialEq, Eq)]
175pub struct ElectLeadersTopicResultV0 {
176    pub name: String,
177    pub partitions: Vec<ElectLeadersPartitionResultV0>,
178}
179
180#[derive(Debug, Clone, PartialEq, Eq)]
181pub struct ElectLeadersPartitionResultV0 {
182    pub partition_index: i32,
183    pub error_code: i16,
184    pub error_message: Option<String>,
185}
186
187fn encode_legacy_topics(
188    encoder: &mut Encoder,
189    topics: Option<&[ElectLeadersTopicV0]>,
190) -> Result<()> {
191    encoder.write_array(topics, |encoder, topic| {
192        encoder.write_string(&topic.name)?;
193        encoder.write_array(Some(&topic.partitions), |encoder, partition| {
194            encoder.write_i32(*partition);
195            Ok(())
196        })
197    })
198}
199
200fn decode_legacy_topic_result(decoder: &mut Decoder<'_>) -> Result<ElectLeadersTopicResultV0> {
201    let name = decoder.read_string()?;
202    let partitions = decoder
203        .read_array("elect leaders partition results", |decoder| {
204            Ok(ElectLeadersPartitionResultV0 {
205                partition_index: decoder.read_i32()?,
206                error_code: decoder.read_i16()?,
207                error_message: decoder.read_nullable_string()?,
208            })
209        })?
210        .unwrap_or_default();
211    Ok(ElectLeadersTopicResultV0 { name, partitions })
212}
213
214#[cfg(test)]
215#[allow(clippy::unwrap_used)]
216mod tests {
217    use super::{
218        ElectLeadersRequestV0, ElectLeadersRequestV1, ElectLeadersRequestV2,
219        ElectLeadersResponseV1, ElectLeadersResponseV2, ElectLeadersTopicV0, API_KEY,
220    };
221    use crate::codec::{Decoder, Encoder};
222
223    #[test]
224    fn encodes_elect_leaders_v0_with_all_topics() {
225        let request = ElectLeadersRequestV0 {
226            correlation_id: 43,
227            client_id: Some("kafrust".to_owned()),
228            topics: None,
229            timeout_ms: 30_000,
230        };
231
232        let bytes = request.encode().unwrap();
233        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 0]);
234        assert_eq!(&bytes[4..8], &[0, 0, 0, 43]);
235        assert_eq!(&bytes[17..21], &[255, 255, 255, 255]);
236        assert_eq!(&bytes[21..25], &30_000_i32.to_be_bytes());
237    }
238
239    #[test]
240    fn encodes_elect_leaders_v1_with_partition_filter() {
241        let request = ElectLeadersRequestV1 {
242            correlation_id: 44,
243            client_id: None,
244            election_type: 1,
245            topics: Some(vec![ElectLeadersTopicV0 {
246                name: "orders".to_owned(),
247                partitions: vec![0, 2],
248            }]),
249            timeout_ms: 10_000,
250        };
251
252        let bytes = request.encode().unwrap();
253        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 1]);
254        assert_eq!(&bytes[4..8], &[0, 0, 0, 44]);
255        assert_eq!(bytes[10], 1);
256        assert!(bytes
257            .windows(6)
258            .any(|window| window == [0, 6, b'o', b'r', b'd', b'e']));
259        assert!(bytes.ends_with(&10_000_i32.to_be_bytes()));
260    }
261
262    #[test]
263    fn encodes_elect_leaders_v2_with_flexible_fields() {
264        let request = ElectLeadersRequestV2 {
265            correlation_id: 45,
266            client_id: None,
267            election_type: 0,
268            topics: Some(vec![ElectLeadersTopicV0 {
269                name: "orders".to_owned(),
270                partitions: vec![1],
271            }]),
272            timeout_ms: 5_000,
273        };
274
275        let bytes = request.encode().unwrap();
276        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 2]);
277        assert_eq!(&bytes[4..8], &[0, 0, 0, 45]);
278        assert_eq!(bytes[11], 0);
279        assert!(bytes.ends_with(&[0]));
280    }
281
282    #[test]
283    fn decodes_elect_leaders_v1_response() {
284        let mut bytes = Encoder::new();
285        bytes.write_i32(8);
286        bytes.write_i16(0);
287        bytes
288            .write_array(
289                Some(&[ElectLeadersTopicV0 {
290                    name: "orders".to_owned(),
291                    partitions: vec![0],
292                }]),
293                |encoder, topic| {
294                    encoder.write_string(&topic.name)?;
295                    encoder.write_array(Some(&[0_i32]), |encoder, partition| {
296                        encoder.write_i32(*partition);
297                        encoder.write_i16(0);
298                        encoder.write_nullable_string(Some("ok"))
299                    })
300                },
301            )
302            .unwrap();
303        let encoded = bytes.into_bytes();
304        let mut decoder = Decoder::new(&encoded);
305        let response = ElectLeadersResponseV1::decode_body(&mut decoder).unwrap();
306
307        assert_eq!(response.throttle_time_ms, 8);
308        assert_eq!(response.results[0].name, "orders");
309        assert_eq!(response.results[0].partitions[0].partition_index, 0);
310        assert_eq!(
311            response.results[0].partitions[0].error_message.as_deref(),
312            Some("ok")
313        );
314        assert!(decoder.is_empty());
315    }
316
317    #[test]
318    fn decodes_elect_leaders_v2_response_with_tagged_fields() {
319        let mut bytes = Encoder::new();
320        bytes.write_i32(4);
321        bytes.write_i16(0);
322        bytes.write_unsigned_varint(2);
323        bytes.write_compact_string("orders").unwrap();
324        bytes.write_unsigned_varint(2);
325        bytes.write_i32(1);
326        bytes.write_i16(0);
327        bytes.write_compact_nullable_string(None).unwrap();
328        bytes.write_empty_tagged_fields();
329        bytes.write_empty_tagged_fields();
330        bytes.write_empty_tagged_fields();
331        let encoded = bytes.into_bytes();
332        let mut decoder = Decoder::new(&encoded);
333        let response = ElectLeadersResponseV2::decode_body(&mut decoder).unwrap();
334
335        assert_eq!(response.throttle_time_ms, 4);
336        assert_eq!(response.results[0].partitions[0].partition_index, 1);
337        assert!(decoder.is_empty());
338    }
339}