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 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 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 pub fn decode_body_v1(decoder: &mut Decoder<'_>) -> Result<Self> {
66 decode_body(decoder, false)
67 }
68
69 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}