Skip to main content

kafrust_protocol/api/
list_partition_reassignments.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 46;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct ListPartitionReassignmentsRequestV0 {
9    pub correlation_id: i32,
10    pub client_id: Option<String>,
11    pub timeout_ms: i32,
12    pub topics: Option<Vec<ListPartitionReassignmentsTopicV0>>,
13}
14
15impl ListPartitionReassignmentsRequestV0 {
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_v2(&mut encoder)?;
25        encoder.write_i32(self.timeout_ms);
26        encoder.write_compact_array(self.topics.as_deref(), |encoder, topic| {
27            encoder.write_compact_string(&topic.name)?;
28            encoder.write_compact_array(Some(&topic.partition_indexes), |encoder, index| {
29                encoder.write_i32(*index);
30                Ok(())
31            })?;
32            encoder.write_empty_tagged_fields();
33            Ok(())
34        })?;
35        encoder.write_empty_tagged_fields();
36        Ok(encoder.into_bytes())
37    }
38}
39
40#[derive(Debug, Clone, PartialEq, Eq)]
41pub struct ListPartitionReassignmentsTopicV0 {
42    pub name: String,
43    pub partition_indexes: Vec<i32>,
44}
45
46#[derive(Debug, Clone, PartialEq, Eq)]
47pub struct ListPartitionReassignmentsResponseV0 {
48    pub throttle_time_ms: i32,
49    pub error_code: i16,
50    pub error_message: Option<String>,
51    pub topics: Vec<ListPartitionReassignmentsTopicResponseV0>,
52}
53
54impl ListPartitionReassignmentsResponseV0 {
55    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
56        let throttle_time_ms = decoder.read_i32()?;
57        let error_code = decoder.read_i16()?;
58        let error_message = decoder.read_compact_nullable_string()?;
59        let topics = decoder
60            .read_compact_array("list partition reassignment topics", |decoder| {
61                let name = decoder.read_compact_string()?;
62                let partitions = decoder
63                    .read_compact_array("list partition reassignment partitions", |decoder| {
64                        let partition_index = decoder.read_i32()?;
65                        let replicas = decoder
66                            .read_array("partition reassignment replicas", |decoder| {
67                                decoder.read_i32()
68                            })?
69                            .unwrap_or_default();
70                        let adding_replicas = decoder
71                            .read_array("partition reassignment adding replicas", |decoder| {
72                                decoder.read_i32()
73                            })?
74                            .unwrap_or_default();
75                        let removing_replicas = decoder
76                            .read_array("partition reassignment removing replicas", |decoder| {
77                                decoder.read_i32()
78                            })?
79                            .unwrap_or_default();
80                        decoder.read_tagged_fields()?;
81                        Ok(ListPartitionReassignmentsPartitionResponseV0 {
82                            partition_index,
83                            replicas,
84                            adding_replicas,
85                            removing_replicas,
86                        })
87                    })?
88                    .unwrap_or_default();
89                decoder.read_tagged_fields()?;
90                Ok(ListPartitionReassignmentsTopicResponseV0 { name, partitions })
91            })?
92            .unwrap_or_default();
93        decoder.read_tagged_fields()?;
94        Ok(Self {
95            throttle_time_ms,
96            error_code,
97            error_message,
98            topics,
99        })
100    }
101}
102
103#[derive(Debug, Clone, PartialEq, Eq)]
104pub struct ListPartitionReassignmentsTopicResponseV0 {
105    pub name: String,
106    pub partitions: Vec<ListPartitionReassignmentsPartitionResponseV0>,
107}
108
109#[derive(Debug, Clone, PartialEq, Eq)]
110pub struct ListPartitionReassignmentsPartitionResponseV0 {
111    pub partition_index: i32,
112    pub replicas: Vec<i32>,
113    pub adding_replicas: Vec<i32>,
114    pub removing_replicas: Vec<i32>,
115}
116
117#[cfg(test)]
118#[allow(clippy::unwrap_used)]
119mod tests {
120    use super::{
121        ListPartitionReassignmentsRequestV0, ListPartitionReassignmentsResponseV0,
122        ListPartitionReassignmentsTopicV0, API_KEY,
123    };
124    use crate::codec::{Decoder, Encoder};
125
126    #[test]
127    fn encodes_list_partition_reassignments_v0_request_with_nullable_topics() {
128        let request = ListPartitionReassignmentsRequestV0 {
129            correlation_id: 13,
130            client_id: None,
131            timeout_ms: 10_000,
132            topics: Some(vec![ListPartitionReassignmentsTopicV0 {
133                name: "orders".to_owned(),
134                partition_indexes: vec![0, 2],
135            }]),
136        };
137
138        let bytes = request.encode().unwrap();
139        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 0]);
140        assert_eq!(&bytes[4..8], &[0, 0, 0, 13]);
141        assert_eq!(bytes[bytes.len() - 1], 0);
142    }
143
144    #[test]
145    fn decodes_list_partition_reassignments_v0_response() {
146        let mut bytes = Encoder::new();
147        bytes.write_i32(9);
148        bytes.write_i16(0);
149        bytes.write_compact_nullable_string(None).unwrap();
150        bytes.write_unsigned_varint(2);
151        bytes.write_compact_string("orders").unwrap();
152        bytes.write_unsigned_varint(2);
153        bytes.write_i32(2);
154        bytes
155            .write_array(Some(&[1, 2, 3]), |encoder, value| {
156                encoder.write_i32(*value);
157                Ok(())
158            })
159            .unwrap();
160        bytes
161            .write_array(Some(&[3]), |encoder, value| {
162                encoder.write_i32(*value);
163                Ok(())
164            })
165            .unwrap();
166        bytes
167            .write_array(Some(&[1]), |encoder, value| {
168                encoder.write_i32(*value);
169                Ok(())
170            })
171            .unwrap();
172        bytes.write_empty_tagged_fields();
173        bytes.write_empty_tagged_fields();
174        bytes.write_empty_tagged_fields();
175        let bytes = bytes.into_bytes();
176        let mut decoder = Decoder::new(&bytes);
177
178        let response = ListPartitionReassignmentsResponseV0::decode_body(&mut decoder).unwrap();
179
180        assert_eq!(response.throttle_time_ms, 9);
181        assert_eq!(response.topics[0].name, "orders");
182        assert_eq!(response.topics[0].partitions[0].replicas, [1, 2, 3]);
183        assert_eq!(response.topics[0].partitions[0].adding_replicas, [3]);
184        assert_eq!(response.topics[0].partitions[0].removing_replicas, [1]);
185        assert!(decoder.is_empty());
186    }
187}