Skip to main content

kafrust_protocol/api/
alter_replica_log_dirs.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 34;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct AlterReplicaLogDirsRequest {
9    pub correlation_id: i32,
10    pub client_id: Option<String>,
11    pub dirs: Vec<AlterReplicaLogDir>,
12}
13
14impl AlterReplicaLogDirsRequest {
15    /// Encodes the non-flexible v1 request schema.
16    pub fn encode_v1(&self) -> Result<Vec<u8>> {
17        let mut encoder = Encoder::new();
18        RequestHeader {
19            api_key: API_KEY,
20            api_version: 1,
21            correlation_id: self.correlation_id,
22            client_id: self.client_id.clone(),
23        }
24        .encode_v1(&mut encoder)?;
25        encode_legacy_dirs(&mut encoder, &self.dirs)?;
26        Ok(encoder.into_bytes())
27    }
28
29    /// Encodes a flexible v2 or newer request schema.
30    pub fn encode_v2(&self, api_version: i16) -> Result<Vec<u8>> {
31        let mut encoder = Encoder::new();
32        RequestHeader {
33            api_key: API_KEY,
34            api_version,
35            correlation_id: self.correlation_id,
36            client_id: self.client_id.clone(),
37        }
38        .encode_v2(&mut encoder)?;
39        encoder.write_compact_array(Some(&self.dirs), encode_flexible_dir)?;
40        encoder.write_empty_tagged_fields();
41        Ok(encoder.into_bytes())
42    }
43}
44
45#[derive(Debug, Clone, PartialEq, Eq)]
46pub struct AlterReplicaLogDir {
47    pub path: String,
48    pub topics: Vec<AlterReplicaLogDirTopic>,
49}
50
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct AlterReplicaLogDirTopic {
53    pub name: String,
54    pub partitions: Vec<i32>,
55}
56
57#[derive(Debug, Clone, PartialEq, Eq)]
58pub struct AlterReplicaLogDirsResponse {
59    pub throttle_time_ms: i32,
60    pub results: Vec<AlterReplicaLogDirTopicResult>,
61}
62
63impl AlterReplicaLogDirsResponse {
64    /// Decodes the non-flexible v1 response schema.
65    pub fn decode_body_v1(decoder: &mut Decoder<'_>) -> Result<Self> {
66        decode_body(decoder, false)
67    }
68
69    /// Decodes the flexible v2 response schema.
70    pub fn decode_body_v2(decoder: &mut Decoder<'_>) -> Result<Self> {
71        decode_body(decoder, true)
72    }
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub struct AlterReplicaLogDirTopicResult {
77    pub name: String,
78    pub partitions: Vec<AlterReplicaLogDirPartitionResult>,
79}
80
81#[derive(Debug, Clone, PartialEq, Eq)]
82pub struct AlterReplicaLogDirPartitionResult {
83    pub partition_index: i32,
84    pub error_code: i16,
85}
86
87fn encode_legacy_dirs(encoder: &mut Encoder, dirs: &[AlterReplicaLogDir]) -> Result<()> {
88    encoder.write_array(Some(dirs), |encoder, dir| {
89        encoder.write_string(&dir.path)?;
90        encoder.write_array(Some(&dir.topics), |encoder, topic| {
91            encoder.write_string(&topic.name)?;
92            encoder.write_array(Some(&topic.partitions), |encoder, partition| {
93                encoder.write_i32(*partition);
94                Ok(())
95            })
96        })
97    })
98}
99
100fn encode_flexible_dir(encoder: &mut Encoder, dir: &AlterReplicaLogDir) -> Result<()> {
101    encoder.write_compact_string(&dir.path)?;
102    encoder.write_compact_array(Some(&dir.topics), |encoder, topic| {
103        encoder.write_compact_string(&topic.name)?;
104        encoder.write_compact_array(Some(&topic.partitions), |encoder, partition| {
105            encoder.write_i32(*partition);
106            Ok(())
107        })?;
108        encoder.write_empty_tagged_fields();
109        Ok(())
110    })?;
111    encoder.write_empty_tagged_fields();
112    Ok(())
113}
114
115fn decode_body(decoder: &mut Decoder<'_>, flexible: bool) -> Result<AlterReplicaLogDirsResponse> {
116    let throttle_time_ms = decoder.read_i32()?;
117    let results = if flexible {
118        decoder
119            .read_compact_array(
120                "alter replica log dirs results",
121                decode_flexible_topic_result,
122            )?
123            .unwrap_or_default()
124    } else {
125        decoder
126            .read_array("alter replica log dirs results", decode_legacy_topic_result)?
127            .unwrap_or_default()
128    };
129    if flexible {
130        decoder.read_tagged_fields()?;
131    }
132    Ok(AlterReplicaLogDirsResponse {
133        throttle_time_ms,
134        results,
135    })
136}
137
138fn decode_legacy_topic_result(decoder: &mut Decoder<'_>) -> Result<AlterReplicaLogDirTopicResult> {
139    Ok(AlterReplicaLogDirTopicResult {
140        name: decoder.read_string()?,
141        partitions: decoder
142            .read_array("alter replica log dirs partitions", decode_partition_result)?
143            .unwrap_or_default(),
144    })
145}
146
147fn decode_flexible_topic_result(
148    decoder: &mut Decoder<'_>,
149) -> Result<AlterReplicaLogDirTopicResult> {
150    let result = AlterReplicaLogDirTopicResult {
151        name: decoder.read_compact_string()?,
152        partitions: decoder
153            .read_compact_array(
154                "alter replica log dirs partitions",
155                decode_flexible_partition_result,
156            )?
157            .unwrap_or_default(),
158    };
159    decoder.read_tagged_fields()?;
160    Ok(result)
161}
162
163fn decode_flexible_partition_result(
164    decoder: &mut Decoder<'_>,
165) -> Result<AlterReplicaLogDirPartitionResult> {
166    let result = decode_partition_result(decoder)?;
167    decoder.read_tagged_fields()?;
168    Ok(result)
169}
170
171fn decode_partition_result(decoder: &mut Decoder<'_>) -> Result<AlterReplicaLogDirPartitionResult> {
172    Ok(AlterReplicaLogDirPartitionResult {
173        partition_index: decoder.read_i32()?,
174        error_code: decoder.read_i16()?,
175    })
176}
177
178#[cfg(test)]
179#[allow(clippy::unwrap_used)]
180mod tests {
181    use super::{
182        AlterReplicaLogDir, AlterReplicaLogDirPartitionResult, AlterReplicaLogDirTopic,
183        AlterReplicaLogDirTopicResult, AlterReplicaLogDirsRequest, AlterReplicaLogDirsResponse,
184        API_KEY,
185    };
186    use crate::codec::{Decoder, Encoder};
187
188    #[test]
189    fn encodes_alter_replica_log_dirs_v1() {
190        let request = AlterReplicaLogDirsRequest {
191            correlation_id: 34,
192            client_id: Some("kafrust".to_owned()),
193            dirs: vec![AlterReplicaLogDir {
194                path: "/var/lib/kafka-2".to_owned(),
195                topics: vec![AlterReplicaLogDirTopic {
196                    name: "orders".to_owned(),
197                    partitions: vec![0, 2],
198                }],
199            }],
200        };
201
202        let bytes = request.encode_v1().unwrap();
203        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 1]);
204        assert_eq!(&bytes[4..8], &[0, 0, 0, 34]);
205        assert!(bytes
206            .windows(8)
207            .any(|window| { window == [0, 6, b'o', b'r', b'd', b'e', b'r', b's'] }));
208        assert!(bytes.windows(4).any(|window| window == [0, 0, 0, 2]));
209    }
210
211    #[test]
212    fn encodes_alter_replica_log_dirs_v2_with_tagged_fields() {
213        let request = AlterReplicaLogDirsRequest {
214            correlation_id: 35,
215            client_id: None,
216            dirs: vec![AlterReplicaLogDir {
217                path: "/var/lib/kafka-2".to_owned(),
218                topics: vec![AlterReplicaLogDirTopic {
219                    name: "orders".to_owned(),
220                    partitions: vec![1],
221                }],
222            }],
223        };
224
225        let bytes = request.encode_v2(2).unwrap();
226        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 2]);
227        assert_eq!(&bytes[4..8], &[0, 0, 0, 35]);
228        assert!(bytes.ends_with(&[0]));
229    }
230
231    #[test]
232    fn decodes_alter_replica_log_dirs_v1_response() {
233        let mut bytes = Encoder::new();
234        bytes.write_i32(7);
235        bytes
236            .write_array(
237                Some(&[AlterReplicaLogDirTopicResult {
238                    name: "orders".to_owned(),
239                    partitions: vec![AlterReplicaLogDirPartitionResult {
240                        partition_index: 0,
241                        error_code: 0,
242                    }],
243                }]),
244                |encoder, topic| {
245                    encoder.write_string(&topic.name)?;
246                    encoder.write_array(Some(&topic.partitions), |encoder, partition| {
247                        encoder.write_i32(partition.partition_index);
248                        encoder.write_i16(partition.error_code);
249                        Ok(())
250                    })
251                },
252            )
253            .unwrap();
254        let encoded = bytes.into_bytes();
255        let mut decoder = Decoder::new(&encoded);
256        let response = AlterReplicaLogDirsResponse::decode_body_v1(&mut decoder).unwrap();
257
258        assert_eq!(response.throttle_time_ms, 7);
259        assert_eq!(response.results[0].name, "orders");
260        assert_eq!(response.results[0].partitions[0].partition_index, 0);
261        assert_eq!(response.results[0].partitions[0].error_code, 0);
262        assert!(decoder.is_empty());
263    }
264
265    #[test]
266    fn decodes_alter_replica_log_dirs_v2_response_with_tagged_fields() {
267        let mut bytes = Encoder::new();
268        bytes.write_i32(8);
269        bytes.write_unsigned_varint(2);
270        bytes.write_compact_string("orders").unwrap();
271        bytes.write_unsigned_varint(2);
272        bytes.write_i32(1);
273        bytes.write_i16(0);
274        bytes.write_empty_tagged_fields();
275        bytes.write_empty_tagged_fields();
276        bytes.write_empty_tagged_fields();
277        let encoded = bytes.into_bytes();
278        let mut decoder = Decoder::new(&encoded);
279        let response = AlterReplicaLogDirsResponse::decode_body_v2(&mut decoder).unwrap();
280
281        assert_eq!(response.throttle_time_ms, 8);
282        assert_eq!(response.results[0].partitions[0].partition_index, 1);
283        assert!(decoder.is_empty());
284    }
285}