Skip to main content

kafrust_protocol/api/
raft_voter.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5/// Kafka AddRaftVoter API key.
6pub const ADD_RAFT_VOTER_API_KEY: i16 = 80;
7/// Kafka RemoveRaftVoter API key.
8pub const REMOVE_RAFT_VOTER_API_KEY: i16 = 81;
9
10/// One listener endpoint supplied when adding a KRaft voter.
11#[derive(Debug, Clone, PartialEq, Eq)]
12pub struct RaftVoterListener {
13    pub name: String,
14    pub host: String,
15    pub port: u16,
16}
17
18/// AddRaftVoter v0 request.
19#[derive(Debug, Clone, PartialEq, Eq)]
20pub struct AddRaftVoterRequestV0 {
21    pub correlation_id: i32,
22    pub client_id: Option<String>,
23    pub cluster_id: Option<String>,
24    pub timeout_ms: i32,
25    pub voter_id: i32,
26    pub voter_directory_id: [u8; 16],
27    pub listeners: Vec<RaftVoterListener>,
28}
29
30impl AddRaftVoterRequestV0 {
31    /// Encodes the flexible v0 request header and body.
32    pub fn encode(&self) -> Result<Vec<u8>> {
33        let mut encoder = Encoder::new();
34        RequestHeader {
35            api_key: ADD_RAFT_VOTER_API_KEY,
36            api_version: 0,
37            correlation_id: self.correlation_id,
38            client_id: self.client_id.clone(),
39        }
40        .encode_v2(&mut encoder)?;
41        encode_body(
42            &mut encoder,
43            self.cluster_id.as_deref(),
44            self.timeout_ms,
45            self.voter_id,
46            &self.voter_directory_id,
47            &self.listeners,
48        )?;
49        encoder.write_empty_tagged_fields();
50        Ok(encoder.into_bytes())
51    }
52}
53
54/// AddRaftVoter v1 request.
55#[derive(Debug, Clone, PartialEq, Eq)]
56pub struct AddRaftVoterRequestV1 {
57    pub correlation_id: i32,
58    pub client_id: Option<String>,
59    pub cluster_id: Option<String>,
60    pub timeout_ms: i32,
61    pub voter_id: i32,
62    pub voter_directory_id: [u8; 16],
63    pub listeners: Vec<RaftVoterListener>,
64    pub ack_when_committed: bool,
65}
66
67impl AddRaftVoterRequestV1 {
68    /// Encodes the flexible v1 request header and body.
69    pub fn encode(&self) -> Result<Vec<u8>> {
70        let mut encoder = Encoder::new();
71        RequestHeader {
72            api_key: ADD_RAFT_VOTER_API_KEY,
73            api_version: 1,
74            correlation_id: self.correlation_id,
75            client_id: self.client_id.clone(),
76        }
77        .encode_v2(&mut encoder)?;
78        encode_body(
79            &mut encoder,
80            self.cluster_id.as_deref(),
81            self.timeout_ms,
82            self.voter_id,
83            &self.voter_directory_id,
84            &self.listeners,
85        )?;
86        encoder.write_bool(self.ack_when_committed);
87        encoder.write_empty_tagged_fields();
88        Ok(encoder.into_bytes())
89    }
90}
91
92/// RemoveRaftVoter v0 request.
93#[derive(Debug, Clone, PartialEq, Eq)]
94pub struct RemoveRaftVoterRequestV0 {
95    pub correlation_id: i32,
96    pub client_id: Option<String>,
97    pub cluster_id: Option<String>,
98    pub voter_id: i32,
99    pub voter_directory_id: [u8; 16],
100}
101
102impl RemoveRaftVoterRequestV0 {
103    /// Encodes the flexible v0 request header and body.
104    pub fn encode(&self) -> Result<Vec<u8>> {
105        let mut encoder = Encoder::new();
106        RequestHeader {
107            api_key: REMOVE_RAFT_VOTER_API_KEY,
108            api_version: 0,
109            correlation_id: self.correlation_id,
110            client_id: self.client_id.clone(),
111        }
112        .encode_v2(&mut encoder)?;
113        encoder.write_compact_nullable_string(self.cluster_id.as_deref())?;
114        encoder.write_i32(self.voter_id);
115        encoder.write_uuid(&self.voter_directory_id);
116        encoder.write_empty_tagged_fields();
117        Ok(encoder.into_bytes())
118    }
119}
120
121/// AddRaftVoter response shared by v0 and v1.
122#[derive(Debug, Clone, PartialEq, Eq)]
123pub struct AddRaftVoterResponse {
124    pub throttle_time_ms: i32,
125    pub error_code: i16,
126    pub error_message: Option<String>,
127}
128
129impl AddRaftVoterResponse {
130    /// Decodes a flexible AddRaftVoter response body.
131    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
132        let response = Self {
133            throttle_time_ms: decoder.read_i32()?,
134            error_code: decoder.read_i16()?,
135            error_message: decoder.read_compact_nullable_string()?,
136        };
137        decoder.read_tagged_fields()?;
138        Ok(response)
139    }
140}
141
142/// RemoveRaftVoter response.
143pub type RemoveRaftVoterResponse = AddRaftVoterResponse;
144
145fn encode_body(
146    encoder: &mut Encoder,
147    cluster_id: Option<&str>,
148    timeout_ms: i32,
149    voter_id: i32,
150    voter_directory_id: &[u8; 16],
151    listeners: &[RaftVoterListener],
152) -> Result<()> {
153    encoder.write_compact_nullable_string(cluster_id)?;
154    encoder.write_i32(timeout_ms);
155    encoder.write_i32(voter_id);
156    encoder.write_uuid(voter_directory_id);
157    encoder.write_compact_array(Some(listeners), |encoder, listener| {
158        encoder.write_compact_string(&listener.name)?;
159        encoder.write_compact_string(&listener.host)?;
160        encoder.write_i16(listener.port as i16);
161        encoder.write_empty_tagged_fields();
162        Ok(())
163    })?;
164    Ok(())
165}
166
167#[cfg(test)]
168#[allow(clippy::unwrap_used)]
169mod tests {
170    use super::{
171        AddRaftVoterRequestV0, AddRaftVoterRequestV1, AddRaftVoterResponse, RaftVoterListener,
172        RemoveRaftVoterRequestV0, ADD_RAFT_VOTER_API_KEY, REMOVE_RAFT_VOTER_API_KEY,
173    };
174    use crate::codec::{Decoder, Encoder};
175
176    fn listener() -> RaftVoterListener {
177        RaftVoterListener {
178            name: "CONTROLLER".to_owned(),
179            host: "controller".to_owned(),
180            port: 9093,
181        }
182    }
183
184    #[test]
185    fn encodes_add_raft_voter_v0_wire_shape() {
186        let request = AddRaftVoterRequestV0 {
187            correlation_id: 7,
188            client_id: Some("kafrust".to_owned()),
189            cluster_id: Some("cluster".to_owned()),
190            timeout_ms: 30_000,
191            voter_id: 4,
192            voter_directory_id: [9; 16],
193            listeners: vec![listener()],
194        };
195        let bytes = request.encode().unwrap();
196        let mut decoder = Decoder::new(&bytes);
197        assert_eq!(decoder.read_i16().unwrap(), ADD_RAFT_VOTER_API_KEY);
198        assert_eq!(decoder.read_i16().unwrap(), 0);
199        assert_eq!(decoder.read_i32().unwrap(), 7);
200        assert_eq!(
201            decoder.read_nullable_string().unwrap().as_deref(),
202            Some("kafrust")
203        );
204        decoder.read_tagged_fields().unwrap();
205        assert_eq!(
206            decoder.read_compact_nullable_string().unwrap().as_deref(),
207            Some("cluster")
208        );
209        assert_eq!(decoder.read_i32().unwrap(), 30_000);
210        assert_eq!(decoder.read_i32().unwrap(), 4);
211        assert_eq!(decoder.read_uuid().unwrap(), [9; 16]);
212        let listeners = decoder
213            .read_compact_array("listeners", |decoder| {
214                let listener = (
215                    decoder.read_compact_string()?,
216                    decoder.read_compact_string()?,
217                    decoder.read_i16()? as u16,
218                );
219                decoder.read_tagged_fields()?;
220                Ok(listener)
221            })
222            .unwrap()
223            .unwrap();
224        assert_eq!(
225            listeners,
226            vec![("CONTROLLER".to_owned(), "controller".to_owned(), 9093)]
227        );
228        decoder.read_tagged_fields().unwrap();
229        assert!(decoder.is_empty());
230    }
231
232    #[test]
233    fn encodes_add_raft_voter_v1_ack_flag() {
234        let request = AddRaftVoterRequestV1 {
235            correlation_id: 8,
236            client_id: None,
237            cluster_id: None,
238            timeout_ms: 1,
239            voter_id: 2,
240            voter_directory_id: [0; 16],
241            listeners: Vec::new(),
242            ack_when_committed: true,
243        };
244        let bytes = request.encode().unwrap();
245        assert_eq!(&bytes[..4], &[0, ADD_RAFT_VOTER_API_KEY as u8, 0, 1]);
246        assert_eq!(*bytes.last().unwrap(), 0);
247        assert!(bytes.windows(2).any(|window| window == [1, 0]));
248    }
249
250    #[test]
251    fn encodes_remove_raft_voter_v0_wire_shape() {
252        let request = RemoveRaftVoterRequestV0 {
253            correlation_id: 9,
254            client_id: None,
255            cluster_id: Some("cluster".to_owned()),
256            voter_id: 2,
257            voter_directory_id: [3; 16],
258        };
259        let bytes = request.encode().unwrap();
260        let mut decoder = Decoder::new(&bytes);
261        assert_eq!(decoder.read_i16().unwrap(), REMOVE_RAFT_VOTER_API_KEY);
262        assert_eq!(decoder.read_i16().unwrap(), 0);
263        assert_eq!(decoder.read_i32().unwrap(), 9);
264        assert_eq!(decoder.read_nullable_string().unwrap(), None);
265        decoder.read_tagged_fields().unwrap();
266        assert_eq!(
267            decoder.read_compact_nullable_string().unwrap().as_deref(),
268            Some("cluster")
269        );
270        assert_eq!(decoder.read_i32().unwrap(), 2);
271        assert_eq!(decoder.read_uuid().unwrap(), [3; 16]);
272        decoder.read_tagged_fields().unwrap();
273        assert!(decoder.is_empty());
274    }
275
276    #[test]
277    fn decodes_raft_voter_response() {
278        let mut encoder = Encoder::new();
279        encoder.write_i32(12);
280        encoder.write_i16(0);
281        encoder.write_compact_nullable_string(Some("ok")).unwrap();
282        encoder.write_empty_tagged_fields();
283        let bytes = encoder.into_bytes();
284        let mut decoder = Decoder::new(&bytes);
285        let response = AddRaftVoterResponse::decode_body(&mut decoder).unwrap();
286        assert_eq!(response.throttle_time_ms, 12);
287        assert_eq!(response.error_code, 0);
288        assert_eq!(response.error_message.as_deref(), Some("ok"));
289        assert!(decoder.is_empty());
290    }
291}