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}