Skip to main content

kafrust_protocol/api/
describe_quorum.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5/// Kafka DescribeQuorum API key.
6pub const API_KEY: i16 = 55;
7
8/// One topic and partition selection in a DescribeQuorum request.
9#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct DescribeQuorumTopic {
11    pub name: String,
12    pub partition_indexes: Vec<i32>,
13}
14
15/// DescribeQuorum request shared by versions 0 through 2.
16#[derive(Debug, Clone, PartialEq, Eq)]
17pub struct DescribeQuorumRequest {
18    pub correlation_id: i32,
19    pub client_id: Option<String>,
20    pub topics: Vec<DescribeQuorumTopic>,
21}
22
23impl DescribeQuorumRequest {
24    /// Encodes the flexible request header and body.
25    pub fn encode(&self, api_version: i16) -> Result<Vec<u8>> {
26        let mut encoder = Encoder::new();
27        RequestHeader {
28            api_key: API_KEY,
29            api_version,
30            correlation_id: self.correlation_id,
31            client_id: self.client_id.clone(),
32        }
33        .encode_v2(&mut encoder)?;
34        encoder.write_compact_array(Some(&self.topics), |encoder, topic| {
35            encoder.write_compact_string(&topic.name)?;
36            encoder.write_compact_array(Some(&topic.partition_indexes), |encoder, index| {
37                encoder.write_i32(*index);
38                encoder.write_empty_tagged_fields();
39                Ok(())
40            })?;
41            encoder.write_empty_tagged_fields();
42            Ok(())
43        })?;
44        encoder.write_empty_tagged_fields();
45        Ok(encoder.into_bytes())
46    }
47}
48
49/// Replica state returned by DescribeQuorum.
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub struct DescribeQuorumReplicaState {
52    pub replica_id: i32,
53    pub replica_directory_id: Option<[u8; 16]>,
54    pub log_end_offset: i64,
55    pub last_fetch_timestamp: Option<i64>,
56    pub last_caught_up_timestamp: Option<i64>,
57}
58
59/// One partition returned by DescribeQuorum.
60#[derive(Debug, Clone, PartialEq, Eq)]
61pub struct DescribeQuorumPartition {
62    pub partition_index: i32,
63    pub error_code: i16,
64    pub error_message: Option<String>,
65    pub leader_id: i32,
66    pub leader_epoch: i32,
67    pub high_watermark: i64,
68    pub current_voters: Vec<DescribeQuorumReplicaState>,
69    pub observers: Vec<DescribeQuorumReplicaState>,
70}
71
72/// One topic returned by DescribeQuorum.
73#[derive(Debug, Clone, PartialEq, Eq)]
74pub struct DescribeQuorumTopicResponse {
75    pub name: String,
76    pub partitions: Vec<DescribeQuorumPartition>,
77}
78
79/// One controller listener returned by DescribeQuorum v2.
80#[derive(Debug, Clone, PartialEq, Eq)]
81pub struct DescribeQuorumListener {
82    pub name: String,
83    pub host: String,
84    pub port: u16,
85}
86
87/// One controller node returned by DescribeQuorum v2.
88#[derive(Debug, Clone, PartialEq, Eq)]
89pub struct DescribeQuorumNode {
90    pub node_id: i32,
91    pub listeners: Vec<DescribeQuorumListener>,
92}
93
94/// DescribeQuorum response for versions 0 through 2.
95#[derive(Debug, Clone, PartialEq, Eq)]
96pub struct DescribeQuorumResponse {
97    pub api_version: i16,
98    pub error_code: i16,
99    pub error_message: Option<String>,
100    pub topics: Vec<DescribeQuorumTopicResponse>,
101    pub nodes: Vec<DescribeQuorumNode>,
102}
103
104impl DescribeQuorumResponse {
105    /// Decodes a flexible response body for the negotiated API version.
106    pub fn decode_body(decoder: &mut Decoder<'_>, api_version: i16) -> Result<Self> {
107        let error_code = decoder.read_i16()?;
108        let error_message = if api_version >= 2 {
109            decoder.read_compact_nullable_string()?
110        } else {
111            None
112        };
113        let topics = decoder
114            .read_compact_array("describe quorum topics", |decoder| {
115                let name = decoder.read_compact_string()?;
116                let partitions = decoder
117                    .read_compact_array("describe quorum partitions", |decoder| {
118                        let partition_index = decoder.read_i32()?;
119                        let error_code = decoder.read_i16()?;
120                        let error_message = if api_version >= 2 {
121                            decoder.read_compact_nullable_string()?
122                        } else {
123                            None
124                        };
125                        let leader_id = decoder.read_i32()?;
126                        let leader_epoch = decoder.read_i32()?;
127                        let high_watermark = decoder.read_i64()?;
128                        let current_voters = decoder
129                            .read_compact_array("describe quorum current voters", |decoder| {
130                                decode_replica_state(decoder, api_version)
131                            })?
132                            .unwrap_or_default();
133                        let observers = decoder
134                            .read_compact_array("describe quorum observers", |decoder| {
135                                decode_replica_state(decoder, api_version)
136                            })?
137                            .unwrap_or_default();
138                        decoder.read_tagged_fields()?;
139                        Ok(DescribeQuorumPartition {
140                            partition_index,
141                            error_code,
142                            error_message,
143                            leader_id,
144                            leader_epoch,
145                            high_watermark,
146                            current_voters,
147                            observers,
148                        })
149                    })?
150                    .unwrap_or_default();
151                decoder.read_tagged_fields()?;
152                Ok(DescribeQuorumTopicResponse { name, partitions })
153            })?
154            .unwrap_or_default();
155        let nodes = if api_version >= 2 {
156            decoder
157                .read_compact_array("describe quorum nodes", |decoder| {
158                    let node_id = decoder.read_i32()?;
159                    let listeners = decoder
160                        .read_compact_array("describe quorum listeners", |decoder| {
161                            let name = decoder.read_compact_string()?;
162                            let host = decoder.read_compact_string()?;
163                            let port = decoder.read_i16()? as u16;
164                            decoder.read_tagged_fields()?;
165                            Ok(DescribeQuorumListener { name, host, port })
166                        })?
167                        .unwrap_or_default();
168                    decoder.read_tagged_fields()?;
169                    Ok(DescribeQuorumNode { node_id, listeners })
170                })?
171                .unwrap_or_default()
172        } else {
173            Vec::new()
174        };
175        decoder.read_tagged_fields()?;
176        Ok(Self {
177            api_version,
178            error_code,
179            error_message,
180            topics,
181            nodes,
182        })
183    }
184}
185
186fn decode_replica_state(
187    decoder: &mut Decoder<'_>,
188    api_version: i16,
189) -> Result<DescribeQuorumReplicaState> {
190    let replica_id = decoder.read_i32()?;
191    let replica_directory_id = if api_version >= 2 {
192        Some(decoder.read_uuid()?)
193    } else {
194        None
195    };
196    let log_end_offset = decoder.read_i64()?;
197    let last_fetch_timestamp = if api_version >= 1 {
198        Some(decoder.read_i64()?)
199    } else {
200        None
201    };
202    let last_caught_up_timestamp = if api_version >= 1 {
203        Some(decoder.read_i64()?)
204    } else {
205        None
206    };
207    decoder.read_tagged_fields()?;
208    Ok(DescribeQuorumReplicaState {
209        replica_id,
210        replica_directory_id,
211        log_end_offset,
212        last_fetch_timestamp,
213        last_caught_up_timestamp,
214    })
215}
216
217#[cfg(test)]
218#[allow(clippy::unwrap_used)]
219mod tests {
220    use super::{DescribeQuorumRequest, DescribeQuorumResponse, DescribeQuorumTopic, API_KEY};
221    use crate::codec::{Decoder, Encoder};
222
223    #[test]
224    fn encodes_describe_quorum_v2_request() {
225        let request = DescribeQuorumRequest {
226            correlation_id: 12,
227            client_id: Some("kafrust".to_owned()),
228            topics: vec![DescribeQuorumTopic {
229                name: "__cluster_metadata".to_owned(),
230                partition_indexes: vec![0],
231            }],
232        };
233        let bytes = request.encode(2).unwrap();
234        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 2]);
235        assert_eq!(&bytes[4..8], &[0, 0, 0, 12]);
236        assert!(bytes
237            .windows("__cluster_metadata".len())
238            .any(|value| value == b"__cluster_metadata"));
239        assert_eq!(*bytes.last().unwrap(), 0);
240    }
241
242    #[test]
243    fn encoded_request_round_trips_against_flexible_v0_wire_shape() {
244        let request = DescribeQuorumRequest {
245            correlation_id: 12,
246            client_id: Some("kafrust".to_owned()),
247            topics: vec![DescribeQuorumTopic {
248                name: "__cluster_metadata".to_owned(),
249                partition_indexes: vec![0],
250            }],
251        };
252        let encoded = request.encode(0).unwrap();
253        let mut decoder = Decoder::new(&encoded);
254
255        assert_eq!(decoder.read_i16().unwrap(), API_KEY);
256        assert_eq!(decoder.read_i16().unwrap(), 0);
257        assert_eq!(decoder.read_i32().unwrap(), 12);
258        assert_eq!(
259            decoder.read_nullable_string().unwrap().as_deref(),
260            Some("kafrust")
261        );
262        decoder.read_tagged_fields().unwrap();
263
264        let topics = decoder
265            .read_compact_array("describe quorum request topics", |decoder| {
266                let name = decoder.read_compact_string()?;
267                let partitions = decoder
268                    .read_compact_array("describe quorum request partitions", |decoder| {
269                        let index = decoder.read_i32()?;
270                        decoder.read_tagged_fields()?;
271                        Ok(index)
272                    })?
273                    .unwrap_or_default();
274                decoder.read_tagged_fields()?;
275                Ok((name, partitions))
276            })
277            .unwrap()
278            .unwrap();
279        decoder.read_tagged_fields().unwrap();
280
281        assert_eq!(topics, vec![("__cluster_metadata".to_owned(), vec![0])]);
282        assert!(decoder.is_empty());
283    }
284
285    #[test]
286    fn decodes_describe_quorum_v2_response_with_nodes_and_replica_state() {
287        let mut bytes = Encoder::new();
288        bytes.write_i16(0);
289        bytes.write_compact_nullable_string(None).unwrap();
290        bytes.write_unsigned_varint(2);
291        bytes.write_compact_string("__cluster_metadata").unwrap();
292        bytes.write_unsigned_varint(2);
293        bytes.write_i32(0);
294        bytes.write_i16(0);
295        bytes.write_compact_nullable_string(None).unwrap();
296        bytes.write_i32(1);
297        bytes.write_i32(4);
298        bytes.write_i64(42);
299        bytes.write_unsigned_varint(2);
300        bytes.write_i32(1);
301        bytes.write_uuid(&[8; 16]);
302        bytes.write_i64(42);
303        bytes.write_i64(100);
304        bytes.write_i64(101);
305        bytes.write_empty_tagged_fields();
306        bytes.write_unsigned_varint(1);
307        bytes.write_empty_tagged_fields();
308        bytes.write_empty_tagged_fields();
309        bytes.write_unsigned_varint(2);
310        bytes.write_i32(1);
311        bytes.write_unsigned_varint(2);
312        bytes.write_compact_string("CONTROLLER").unwrap();
313        bytes.write_compact_string("127.0.0.1").unwrap();
314        bytes.write_i16(9093);
315        bytes.write_empty_tagged_fields();
316        bytes.write_empty_tagged_fields();
317        bytes.write_empty_tagged_fields();
318
319        let encoded = bytes.into_bytes();
320        let mut decoder = Decoder::new(&encoded);
321        let response = DescribeQuorumResponse::decode_body(&mut decoder, 2).unwrap();
322        assert_eq!(response.error_code, 0);
323        assert_eq!(response.topics[0].partitions[0].leader_id, 1);
324        assert_eq!(response.topics[0].partitions[0].high_watermark, 42);
325        assert_eq!(
326            response.topics[0].partitions[0].current_voters[0].replica_directory_id,
327            Some([8; 16])
328        );
329        assert_eq!(response.nodes[0].listeners[0].port, 9093);
330        assert!(decoder.is_empty());
331    }
332}