Skip to main content

kafrust_protocol/api/
describe_topic_partitions.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::{Error, Result};
3use crate::header::RequestHeader;
4
5/// Kafka DescribeTopicPartitions API key.
6pub const API_KEY: i16 = 75;
7
8/// Nullable cursor used to page through topic partition metadata.
9#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct DescribeTopicPartitionsCursorV0 {
11    pub topic_name: String,
12    pub partition_index: i32,
13}
14
15/// One topic requested by DescribeTopicPartitions v0.
16#[derive(Debug, Clone, PartialEq, Eq)]
17pub struct DescribeTopicPartitionsTopicV0 {
18    pub name: String,
19}
20
21/// DescribeTopicPartitions v0 request.
22#[derive(Debug, Clone, PartialEq, Eq)]
23pub struct DescribeTopicPartitionsRequestV0 {
24    pub correlation_id: i32,
25    pub client_id: Option<String>,
26    pub topics: Vec<DescribeTopicPartitionsTopicV0>,
27    pub response_partition_limit: i32,
28    pub cursor: Option<DescribeTopicPartitionsCursorV0>,
29}
30
31impl DescribeTopicPartitionsRequestV0 {
32    /// Encodes the flexible v0 request, including its request header.
33    pub fn encode(&self) -> Result<Vec<u8>> {
34        let mut encoder = Encoder::new();
35        RequestHeader {
36            api_key: API_KEY,
37            api_version: 0,
38            correlation_id: self.correlation_id,
39            client_id: self.client_id.clone(),
40        }
41        .encode_v2(&mut encoder)?;
42        encoder.write_compact_array(Some(&self.topics), |encoder, topic| {
43            encoder.write_compact_string(&topic.name)?;
44            encoder.write_empty_tagged_fields();
45            Ok(())
46        })?;
47        encoder.write_i32(self.response_partition_limit);
48        match &self.cursor {
49            Some(cursor) => {
50                encoder.write_i8(1);
51                encoder.write_compact_string(&cursor.topic_name)?;
52                encoder.write_i32(cursor.partition_index);
53                encoder.write_empty_tagged_fields();
54            }
55            None => encoder.write_i8(-1),
56        }
57        encoder.write_empty_tagged_fields();
58        Ok(encoder.into_bytes())
59    }
60}
61
62/// One partition returned by DescribeTopicPartitions v0.
63#[derive(Debug, Clone, PartialEq, Eq)]
64pub struct DescribeTopicPartitionsPartitionResponseV0 {
65    pub error_code: i16,
66    pub partition_index: i32,
67    pub leader_id: i32,
68    pub leader_epoch: i32,
69    pub replica_nodes: Vec<i32>,
70    pub isr_nodes: Vec<i32>,
71    pub eligible_leader_replicas: Option<Vec<i32>>,
72    pub last_known_elr: Option<Vec<i32>>,
73    pub offline_replicas: Vec<i32>,
74}
75
76/// One topic returned by DescribeTopicPartitions v0.
77#[derive(Debug, Clone, PartialEq, Eq)]
78pub struct DescribeTopicPartitionsTopicResponseV0 {
79    pub error_code: i16,
80    pub name: Option<String>,
81    pub topic_id: [u8; 16],
82    pub is_internal: bool,
83    pub partitions: Vec<DescribeTopicPartitionsPartitionResponseV0>,
84    pub topic_authorized_operations: i32,
85}
86
87/// DescribeTopicPartitions v0 response.
88#[derive(Debug, Clone, PartialEq, Eq)]
89pub struct DescribeTopicPartitionsResponseV0 {
90    pub throttle_time_ms: i32,
91    pub topics: Vec<DescribeTopicPartitionsTopicResponseV0>,
92    pub next_cursor: Option<DescribeTopicPartitionsCursorV0>,
93}
94
95impl DescribeTopicPartitionsResponseV0 {
96    /// Decodes the flexible response body after the response header.
97    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
98        let throttle_time_ms = decoder.read_i32()?;
99        let topics = decoder
100            .read_compact_array("describe topic partitions topics", |decoder| {
101                let error_code = decoder.read_i16()?;
102                let name = decoder.read_compact_nullable_string()?;
103                let topic_id = decoder.read_uuid()?;
104                let is_internal = decoder.read_bool()?;
105                let partitions = decoder
106                    .read_compact_array("describe topic partitions partitions", |decoder| {
107                        let error_code = decoder.read_i16()?;
108                        let partition_index = decoder.read_i32()?;
109                        let leader_id = decoder.read_i32()?;
110                        let leader_epoch = decoder.read_i32()?;
111                        let replica_nodes = decoder
112                            .read_compact_array("describe topic partitions replicas", |decoder| {
113                                decoder.read_i32()
114                            })?
115                            .unwrap_or_default();
116                        let isr_nodes = decoder
117                            .read_compact_array("describe topic partitions isr", |decoder| {
118                                decoder.read_i32()
119                            })?
120                            .unwrap_or_default();
121                        let eligible_leader_replicas = decoder.read_compact_array(
122                            "describe topic partitions eligible leader replicas",
123                            |decoder| decoder.read_i32(),
124                        )?;
125                        let last_known_elr = decoder.read_compact_array(
126                            "describe topic partitions last known elr",
127                            |decoder| decoder.read_i32(),
128                        )?;
129                        let offline_replicas = decoder
130                            .read_compact_array(
131                                "describe topic partitions offline replicas",
132                                |decoder| decoder.read_i32(),
133                            )?
134                            .unwrap_or_default();
135                        decoder.read_tagged_fields()?;
136                        Ok(DescribeTopicPartitionsPartitionResponseV0 {
137                            error_code,
138                            partition_index,
139                            leader_id,
140                            leader_epoch,
141                            replica_nodes,
142                            isr_nodes,
143                            eligible_leader_replicas,
144                            last_known_elr,
145                            offline_replicas,
146                        })
147                    })?
148                    .unwrap_or_default();
149                let topic_authorized_operations = decoder.read_i32()?;
150                decoder.read_tagged_fields()?;
151                Ok(DescribeTopicPartitionsTopicResponseV0 {
152                    error_code,
153                    name,
154                    topic_id,
155                    is_internal,
156                    partitions,
157                    topic_authorized_operations,
158                })
159            })?
160            .unwrap_or_default();
161        let next_cursor = match decoder.read_i8()? {
162            -1 => None,
163            1 => {
164                let topic_name = decoder.read_compact_string()?;
165                let partition_index = decoder.read_i32()?;
166                decoder.read_tagged_fields()?;
167                Some(DescribeTopicPartitionsCursorV0 {
168                    topic_name,
169                    partition_index,
170                })
171            }
172            marker => return Err(Error::InvalidNullableStruct(marker)),
173        };
174        decoder.read_tagged_fields()?;
175        Ok(Self {
176            throttle_time_ms,
177            topics,
178            next_cursor,
179        })
180    }
181}
182
183#[cfg(test)]
184#[allow(clippy::unwrap_used)]
185mod tests {
186    use super::{
187        DescribeTopicPartitionsCursorV0, DescribeTopicPartitionsRequestV0,
188        DescribeTopicPartitionsResponseV0, DescribeTopicPartitionsTopicV0, API_KEY,
189    };
190    use crate::codec::{Decoder, Encoder};
191
192    #[test]
193    fn encodes_describe_topic_partitions_v0_request_with_cursor() {
194        let request = DescribeTopicPartitionsRequestV0 {
195            correlation_id: 31,
196            client_id: Some("kafrust".to_owned()),
197            topics: vec![DescribeTopicPartitionsTopicV0 {
198                name: "orders".to_owned(),
199            }],
200            response_partition_limit: 2000,
201            cursor: Some(DescribeTopicPartitionsCursorV0 {
202                topic_name: "orders".to_owned(),
203                partition_index: 2,
204            }),
205        };
206
207        let bytes = request.encode().unwrap();
208        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 0]);
209        assert_eq!(&bytes[4..8], &[0, 0, 0, 31]);
210        assert!(bytes.windows(6).any(|value| value == b"orders"));
211        assert_eq!(&bytes[bytes.len() - 1..], &[0]);
212    }
213
214    #[test]
215    fn decodes_describe_topic_partitions_v0_response_with_nullable_fields() {
216        let mut bytes = Encoder::new();
217        bytes.write_i32(17);
218        bytes.write_unsigned_varint(2);
219        bytes.write_i16(0);
220        bytes.write_compact_nullable_string(Some("orders")).unwrap();
221        bytes.write_uuid(&[7; 16]);
222        bytes.write_bool(false);
223        bytes.write_unsigned_varint(2);
224        bytes.write_i16(0);
225        bytes.write_i32(0);
226        bytes.write_i32(1);
227        bytes.write_i32(8);
228        bytes.write_unsigned_varint(2);
229        bytes.write_i32(1);
230        bytes.write_unsigned_varint(2);
231        bytes.write_i32(1);
232        bytes.write_unsigned_varint(0);
233        bytes.write_unsigned_varint(0);
234        bytes.write_unsigned_varint(2);
235        bytes.write_i32(2);
236        bytes.write_empty_tagged_fields();
237        bytes.write_i32(-2147483648);
238        bytes.write_empty_tagged_fields();
239        bytes.write_i8(-1);
240        bytes.write_empty_tagged_fields();
241
242        let encoded = bytes.into_bytes();
243        let mut decoder = Decoder::new(&encoded);
244        let response = DescribeTopicPartitionsResponseV0::decode_body(&mut decoder).unwrap();
245        assert_eq!(response.throttle_time_ms, 17);
246        assert_eq!(response.topics.len(), 1);
247        assert_eq!(response.topics[0].name.as_deref(), Some("orders"));
248        assert_eq!(response.topics[0].topic_id, [7; 16]);
249        assert_eq!(response.topics[0].partitions[0].leader_id, 1);
250        assert_eq!(
251            response.topics[0].partitions[0].eligible_leader_replicas,
252            None
253        );
254        assert_eq!(response.next_cursor, None);
255        assert!(decoder.is_empty());
256    }
257}