Skip to main content

kafrust_protocol/api/
update_features.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 57;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct UpdateFeaturesRequestV0 {
9    pub correlation_id: i32,
10    pub client_id: Option<String>,
11    pub timeout_ms: i32,
12    pub updates: Vec<FeatureUpdateV0>,
13}
14
15impl UpdateFeaturesRequestV0 {
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(Some(&self.updates), |encoder, update| {
27            encoder.write_compact_string(&update.feature)?;
28            encoder.write_i16(update.max_version_level);
29            encoder.write_bool(update.allow_downgrade);
30            encoder.write_empty_tagged_fields();
31            Ok(())
32        })?;
33        encoder.write_empty_tagged_fields();
34        Ok(encoder.into_bytes())
35    }
36}
37
38/// Kafka UpdateFeatures v1 request.
39///
40/// Version 1 replaces the v0 boolean downgrade flag with Kafka's three-way
41/// upgrade type and adds validation-only execution.
42#[derive(Debug, Clone, PartialEq, Eq)]
43pub struct UpdateFeaturesRequestV1 {
44    pub correlation_id: i32,
45    pub client_id: Option<String>,
46    pub timeout_ms: i32,
47    pub updates: Vec<FeatureUpdateV1>,
48    pub validate_only: bool,
49}
50
51impl UpdateFeaturesRequestV1 {
52    pub fn encode(&self) -> Result<Vec<u8>> {
53        let mut encoder = Encoder::new();
54        RequestHeader {
55            api_key: API_KEY,
56            api_version: 1,
57            correlation_id: self.correlation_id,
58            client_id: self.client_id.clone(),
59        }
60        .encode_v2(&mut encoder)?;
61        encoder.write_i32(self.timeout_ms);
62        encoder.write_compact_array(Some(&self.updates), |encoder, update| {
63            encoder.write_compact_string(&update.feature)?;
64            encoder.write_i16(update.max_version_level);
65            encoder.write_i8(update.upgrade_type);
66            encoder.write_empty_tagged_fields();
67            Ok(())
68        })?;
69        encoder.write_bool(self.validate_only);
70        encoder.write_empty_tagged_fields();
71        Ok(encoder.into_bytes())
72    }
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub struct FeatureUpdateV0 {
77    pub feature: String,
78    pub max_version_level: i16,
79    pub allow_downgrade: bool,
80}
81
82#[derive(Debug, Clone, PartialEq, Eq)]
83pub struct FeatureUpdateV1 {
84    pub feature: String,
85    pub max_version_level: i16,
86    /// Kafka UpgradeType: 1 upgrade, 2 safe downgrade, 3 unsafe downgrade.
87    pub upgrade_type: i8,
88}
89
90#[derive(Debug, Clone, PartialEq, Eq)]
91pub struct UpdateFeaturesResponseV0 {
92    pub throttle_time_ms: i32,
93    pub error_code: i16,
94    pub error_message: Option<String>,
95    pub results: Vec<FeatureUpdateResultV0>,
96}
97
98impl UpdateFeaturesResponseV0 {
99    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
100        let throttle_time_ms = decoder.read_i32()?;
101        let error_code = decoder.read_i16()?;
102        let error_message = decoder.read_compact_nullable_string()?;
103        let results = decoder
104            .read_compact_array("update features results", |decoder| {
105                let result = FeatureUpdateResultV0 {
106                    feature: decoder.read_compact_string()?,
107                    error_code: decoder.read_i16()?,
108                    error_message: decoder.read_compact_nullable_string()?,
109                };
110                decoder.read_tagged_fields()?;
111                Ok(result)
112            })?
113            .unwrap_or_default();
114        decoder.read_tagged_fields()?;
115        Ok(Self {
116            throttle_time_ms,
117            error_code,
118            error_message,
119            results,
120        })
121    }
122}
123
124/// UpdateFeatures v1 has the same response body as v0.
125pub type UpdateFeaturesResponseV1 = UpdateFeaturesResponseV0;
126
127#[derive(Debug, Clone, PartialEq, Eq)]
128pub struct FeatureUpdateResultV0 {
129    pub feature: String,
130    pub error_code: i16,
131    pub error_message: Option<String>,
132}
133
134#[cfg(test)]
135#[allow(clippy::unwrap_used)]
136mod tests {
137    use super::{
138        FeatureUpdateV0, FeatureUpdateV1, UpdateFeaturesRequestV0, UpdateFeaturesRequestV1,
139        UpdateFeaturesResponseV0, API_KEY,
140    };
141    use crate::codec::Decoder;
142
143    #[test]
144    fn encodes_update_features_v0_request() {
145        let request = UpdateFeaturesRequestV0 {
146            correlation_id: 23,
147            client_id: Some("kafrust".to_owned()),
148            timeout_ms: 60_000,
149            updates: vec![FeatureUpdateV0 {
150                feature: "metadata.version".to_owned(),
151                max_version_level: 21,
152                allow_downgrade: false,
153            }],
154        };
155
156        let bytes = request.encode().unwrap();
157        assert_eq!(&bytes[..2], &API_KEY.to_be_bytes());
158        assert_eq!(&bytes[2..4], &[0, 0]);
159        assert_eq!(&bytes[4..8], &23_i32.to_be_bytes());
160        assert!(bytes.ends_with(&[0]));
161    }
162
163    #[test]
164    fn encodes_update_features_v1_request_with_validation() {
165        let request = UpdateFeaturesRequestV1 {
166            correlation_id: 23,
167            client_id: None,
168            timeout_ms: 60_000,
169            updates: vec![FeatureUpdateV1 {
170                feature: "metadata.version".to_owned(),
171                max_version_level: 21,
172                upgrade_type: 2,
173            }],
174            validate_only: true,
175        };
176
177        let bytes = request.encode().unwrap();
178        assert_eq!(
179            bytes,
180            [
181                0,
182                API_KEY as u8, // api key
183                0,
184                1, // api version
185                0,
186                0,
187                0,
188                23, // correlation id
189                0xff,
190                0xff, // nullable client id
191                0,    // request header tagged fields
192                0,
193                0,
194                234,
195                96, // timeout ms
196                2,  // compact update count
197                17, // compact feature string length
198                b'm',
199                b'e',
200                b't',
201                b'a',
202                b'd',
203                b'a',
204                b't',
205                b'a',
206                b'.',
207                b'v',
208                b'e',
209                b'r',
210                b's',
211                b'i',
212                b'o',
213                b'n',
214                0,
215                21, // max version level
216                2,  // safe downgrade
217                0,  // update tagged fields
218                1,  // validate only
219                0,  // request tagged fields
220            ]
221        );
222    }
223
224    #[test]
225    fn decodes_update_features_v0_response() {
226        let mut bytes = Vec::new();
227        bytes.extend_from_slice(&12_i32.to_be_bytes());
228        bytes.extend_from_slice(&0_i16.to_be_bytes());
229        bytes.extend_from_slice(&[3, b'o', b'k']);
230        bytes.push(2);
231        bytes.push(17);
232        bytes.extend_from_slice(b"metadata.version");
233        bytes.extend_from_slice(&0_i16.to_be_bytes());
234        bytes.extend_from_slice(&[3, b'o', b'k']);
235        bytes.push(0);
236        bytes.push(0);
237
238        let mut decoder = Decoder::new(&bytes);
239        let response = UpdateFeaturesResponseV0::decode_body(&mut decoder).unwrap();
240
241        assert_eq!(response.throttle_time_ms, 12);
242        assert_eq!(response.error_code, 0);
243        assert_eq!(response.error_message.as_deref(), Some("ok"));
244        assert_eq!(response.results.len(), 1);
245        assert_eq!(response.results[0].feature, "metadata.version");
246        assert_eq!(response.results[0].error_message.as_deref(), Some("ok"));
247        assert!(decoder.is_empty());
248    }
249}