Skip to main content

kafrust_protocol/api/
api_versions.rs

1use crate::codec::{Decoder, Encoder, TaggedField};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 18;
6
7#[derive(Debug, Clone, Copy, PartialEq, Eq)]
8pub struct ApiKeyVersion {
9    pub api_key: i16,
10    pub min_version: i16,
11    pub max_version: i16,
12}
13
14/// A feature version range advertised by one broker in ApiVersions v3+.
15#[derive(Debug, Clone, PartialEq, Eq)]
16pub struct SupportedFeature {
17    /// Feature name defined by Kafka's feature registry.
18    pub name: String,
19    /// Minimum version level supported by the broker.
20    pub min_version: i16,
21    /// Maximum version level supported by the broker.
22    pub max_version: i16,
23}
24
25impl SupportedFeature {
26    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
27        let feature = Self {
28            name: decoder.read_compact_string()?,
29            min_version: decoder.read_i16()?,
30            max_version: decoder.read_i16()?,
31        };
32        decoder.read_tagged_fields()?;
33        Ok(feature)
34    }
35}
36
37/// A cluster-wide finalized feature range from ApiVersions v3+.
38#[derive(Debug, Clone, PartialEq, Eq)]
39pub struct FinalizedFeature {
40    /// Feature name defined by Kafka's feature registry.
41    pub name: String,
42    /// Finalized minimum version level.
43    pub min_version_level: i16,
44    /// Finalized maximum version level.
45    pub max_version_level: i16,
46}
47
48impl FinalizedFeature {
49    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
50        let feature = Self {
51            name: decoder.read_compact_string()?,
52            max_version_level: decoder.read_i16()?,
53            min_version_level: decoder.read_i16()?,
54        };
55        decoder.read_tagged_fields()?;
56        Ok(feature)
57    }
58}
59
60/// Common capability lookup implemented by fixed and flexible ApiVersions responses.
61pub trait ApiVersionsLookup {
62    /// Returns the highest broker-supported version not exceeding the client limit.
63    fn highest_supported_version(&self, api_key: i16, max_supported: i16) -> Option<i16>;
64}
65
66impl ApiKeyVersion {
67    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
68        Ok(Self {
69            api_key: decoder.read_i16()?,
70            min_version: decoder.read_i16()?,
71            max_version: decoder.read_i16()?,
72        })
73    }
74
75    fn decode_flexible(decoder: &mut Decoder<'_>) -> Result<Self> {
76        let version = Self::decode(decoder)?;
77        decoder.read_tagged_fields()?;
78        Ok(version)
79    }
80}
81
82#[derive(Debug, Clone, PartialEq, Eq)]
83pub struct ApiVersionsRequestV0 {
84    pub correlation_id: i32,
85    pub client_id: Option<String>,
86}
87
88impl ApiVersionsRequestV0 {
89    pub fn encode(&self) -> Result<Vec<u8>> {
90        let mut encoder = Encoder::new();
91        RequestHeader {
92            api_key: API_KEY,
93            api_version: 0,
94            correlation_id: self.correlation_id,
95            client_id: self.client_id.clone(),
96        }
97        .encode_v1(&mut encoder)?;
98        Ok(encoder.into_bytes())
99    }
100}
101
102/// ApiVersions v3 request with client software identification.
103#[derive(Debug, Clone, PartialEq, Eq)]
104pub struct ApiVersionsRequestV3 {
105    /// Request correlation ID.
106    pub correlation_id: i32,
107    /// Optional Kafka client ID from the request header.
108    pub client_id: Option<String>,
109    /// Client software name reported through KIP-511.
110    pub client_software_name: String,
111    /// Client software version reported through KIP-511.
112    pub client_software_version: String,
113}
114
115impl ApiVersionsRequestV3 {
116    /// Encodes an ApiVersions v3 request frame without the outer frame length.
117    pub fn encode(&self) -> Result<Vec<u8>> {
118        encode_flexible_request(
119            3,
120            self.correlation_id,
121            self.client_id.as_deref(),
122            &self.client_software_name,
123            &self.client_software_version,
124            None,
125            None,
126        )
127    }
128}
129
130/// ApiVersions v4 request.
131///
132/// Version 4 keeps the v3 wire shape while allowing the broker to report
133/// feature minimum versions of zero correctly.
134#[derive(Debug, Clone, PartialEq, Eq)]
135pub struct ApiVersionsRequestV4 {
136    /// Request correlation ID.
137    pub correlation_id: i32,
138    /// Optional Kafka client ID from the request header.
139    pub client_id: Option<String>,
140    /// Client software name reported through KIP-511.
141    pub client_software_name: String,
142    /// Client software version reported through KIP-511.
143    pub client_software_version: String,
144}
145
146impl ApiVersionsRequestV4 {
147    /// Encodes an ApiVersions v4 request frame without the outer frame length.
148    pub fn encode(&self) -> Result<Vec<u8>> {
149        encode_flexible_request(
150            4,
151            self.correlation_id,
152            self.client_id.as_deref(),
153            &self.client_software_name,
154            &self.client_software_version,
155            None,
156            None,
157        )
158    }
159}
160
161/// ApiVersions v5 request with optional cluster and node identity checks.
162#[derive(Debug, Clone, PartialEq, Eq)]
163pub struct ApiVersionsRequestV5 {
164    /// Request correlation ID.
165    pub correlation_id: i32,
166    /// Optional Kafka client ID from the request header.
167    pub client_id: Option<String>,
168    /// Client software name reported through KIP-511.
169    pub client_software_name: String,
170    /// Client software version reported through KIP-511.
171    pub client_software_version: String,
172    /// Expected cluster ID, when the client already knows it.
173    pub cluster_id: Option<String>,
174    /// Expected broker node ID, or `-1` when it is unknown.
175    pub node_id: i32,
176}
177
178impl ApiVersionsRequestV5 {
179    /// Encodes an ApiVersions v5 request frame without the outer frame length.
180    pub fn encode(&self) -> Result<Vec<u8>> {
181        encode_flexible_request(
182            5,
183            self.correlation_id,
184            self.client_id.as_deref(),
185            &self.client_software_name,
186            &self.client_software_version,
187            Some(self.cluster_id.as_deref()),
188            Some(self.node_id),
189        )
190    }
191}
192
193fn encode_flexible_request(
194    api_version: i16,
195    correlation_id: i32,
196    client_id: Option<&str>,
197    client_software_name: &str,
198    client_software_version: &str,
199    cluster_id: Option<Option<&str>>,
200    node_id: Option<i32>,
201) -> Result<Vec<u8>> {
202    let mut encoder = Encoder::new();
203    RequestHeader {
204        api_key: API_KEY,
205        api_version,
206        correlation_id,
207        client_id: client_id.map(str::to_owned),
208    }
209    .encode_v2(&mut encoder)?;
210    encoder.write_compact_string(client_software_name)?;
211    encoder.write_compact_string(client_software_version)?;
212    if let (Some(cluster_id), Some(node_id)) = (cluster_id, node_id) {
213        encoder.write_compact_nullable_string(cluster_id)?;
214        encoder.write_i32(node_id);
215    }
216    encoder.write_empty_tagged_fields();
217    Ok(encoder.into_bytes())
218}
219
220#[derive(Debug, Clone, PartialEq, Eq)]
221pub struct ApiVersionsResponseV0 {
222    pub error_code: i16,
223    pub api_keys: Vec<ApiKeyVersion>,
224}
225
226impl ApiVersionsResponseV0 {
227    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
228        let error_code = decoder.read_i16()?;
229        let api_keys = decoder
230            .read_array("api versions", ApiKeyVersion::decode)?
231            .unwrap_or_default();
232        let response = Self {
233            error_code,
234            api_keys,
235        };
236        decoder.finish()?;
237        Ok(response)
238    }
239
240    pub fn highest_supported_version(&self, api_key: i16, max_supported: i16) -> Option<i16> {
241        self.api_keys
242            .iter()
243            .find(|version| version.api_key == api_key)
244            .and_then(|version| {
245                let selected = version.max_version.min(max_supported);
246                (selected >= version.min_version).then_some(selected)
247            })
248    }
249}
250
251impl ApiVersionsLookup for ApiVersionsResponseV0 {
252    fn highest_supported_version(&self, api_key: i16, max_supported: i16) -> Option<i16> {
253        highest_supported_version(&self.api_keys, api_key, max_supported)
254    }
255}
256
257/// ApiVersions v3 response with flexible encoding and forward-compatible tags.
258#[derive(Debug, Clone, PartialEq, Eq)]
259pub struct ApiVersionsResponseV3 {
260    /// Top-level Kafka error code.
261    pub error_code: i16,
262    /// API version ranges advertised by the broker.
263    pub api_keys: Vec<ApiKeyVersion>,
264    /// Broker throttle duration in milliseconds.
265    pub throttle_time_ms: i32,
266    /// Feature ranges supported by the broker.
267    pub supported_features: Vec<SupportedFeature>,
268    /// Monotonically increasing finalized-feature metadata epoch, or `-1` if unknown.
269    pub finalized_features_epoch: i64,
270    /// Cluster-wide finalized feature ranges when the epoch is known.
271    pub finalized_features: Vec<FinalizedFeature>,
272    /// Whether the broker reports that ZooKeeper migration is ready.
273    pub zk_migration_ready: bool,
274    /// Unknown or future top-level tagged fields preserved for inspection.
275    pub tagged_fields: Vec<TaggedField>,
276}
277
278impl ApiVersionsResponseV3 {
279    /// Decodes an ApiVersions v3 response body after its response header.
280    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
281        let error_code = decoder.read_i16()?;
282        let api_keys = decoder
283            .read_compact_array("api versions", ApiKeyVersion::decode_flexible)?
284            .unwrap_or_default();
285        let throttle_time_ms = decoder.read_i32()?;
286        let tagged_fields = decoder.read_tagged_fields()?;
287        let limits = decoder.limits();
288        let mut supported_features = Vec::new();
289        let mut finalized_features_epoch = -1;
290        let mut finalized_features = Vec::new();
291        let mut zk_migration_ready = false;
292        let mut unknown_tagged_fields = Vec::new();
293        for field in tagged_fields {
294            match field.tag {
295                0 => {
296                    let mut field_decoder = Decoder::with_limits(&field.data, limits);
297                    supported_features = field_decoder
298                        .read_compact_array("supported features", SupportedFeature::decode)?
299                        .unwrap_or_default();
300                }
301                1 => {
302                    let mut field_decoder = Decoder::with_limits(&field.data, limits);
303                    finalized_features_epoch = field_decoder.read_i64()?;
304                }
305                2 => {
306                    let mut field_decoder = Decoder::with_limits(&field.data, limits);
307                    finalized_features = field_decoder
308                        .read_compact_array("finalized features", FinalizedFeature::decode)?
309                        .unwrap_or_default();
310                }
311                3 => {
312                    let mut field_decoder = Decoder::with_limits(&field.data, limits);
313                    zk_migration_ready = field_decoder.read_bool()?;
314                }
315                _ => unknown_tagged_fields.push(field),
316            }
317        }
318        let response = Self {
319            error_code,
320            api_keys,
321            throttle_time_ms,
322            supported_features,
323            finalized_features_epoch,
324            finalized_features,
325            zk_migration_ready,
326            tagged_fields: unknown_tagged_fields,
327        };
328        decoder.finish()?;
329        Ok(response)
330    }
331
332    /// Returns the highest broker-supported version not exceeding the client limit.
333    pub fn highest_supported_version(&self, api_key: i16, max_supported: i16) -> Option<i16> {
334        highest_supported_version(&self.api_keys, api_key, max_supported)
335    }
336}
337
338impl ApiVersionsLookup for ApiVersionsResponseV3 {
339    fn highest_supported_version(&self, api_key: i16, max_supported: i16) -> Option<i16> {
340        highest_supported_version(&self.api_keys, api_key, max_supported)
341    }
342}
343
344/// ApiVersions v4 response. The body schema is identical to v3.
345pub type ApiVersionsResponseV4 = ApiVersionsResponseV3;
346
347/// ApiVersions v5 response. The body schema is identical to v4.
348pub type ApiVersionsResponseV5 = ApiVersionsResponseV3;
349
350/// Kafka's `UNSUPPORTED_VERSION` protocol error code.
351pub const UNSUPPORTED_VERSION_ERROR_CODE: i16 = 35;
352
353fn highest_supported_version(
354    api_keys: &[ApiKeyVersion],
355    api_key: i16,
356    max_supported: i16,
357) -> Option<i16> {
358    api_keys
359        .iter()
360        .find(|version| version.api_key == api_key)
361        .and_then(|version| {
362            let selected = version.max_version.min(max_supported);
363            (selected >= version.min_version).then_some(selected)
364        })
365}
366
367#[cfg(test)]
368#[allow(clippy::unwrap_used)]
369mod tests {
370    use super::{
371        ApiVersionsRequestV0, ApiVersionsRequestV3, ApiVersionsRequestV4, ApiVersionsRequestV5,
372        ApiVersionsResponseV0, ApiVersionsResponseV3, ApiVersionsResponseV4, API_KEY,
373    };
374    use crate::codec::{Decoder, Encoder};
375
376    #[test]
377    fn encodes_api_versions_request_v0() {
378        let request = ApiVersionsRequestV0 {
379            correlation_id: 42,
380            client_id: Some("kafrust".to_owned()),
381        };
382        assert_eq!(
383            request.encode().unwrap(),
384            [
385                0, 18, // api key
386                0, 0, // api version
387                0, 0, 0, 42, // correlation id
388                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't',
389            ]
390        );
391        assert_eq!(API_KEY, 18);
392    }
393
394    #[test]
395    fn decodes_api_versions_response_v0() {
396        let bytes = [
397            0, 0, // error code
398            0, 0, 0, 2, // api key count
399            0, 18, 0, 0, 0, 4, // ApiVersions min/max
400            0, 3, 0, 1, 0, 9, // Metadata min/max
401        ];
402        let mut decoder = Decoder::new(&bytes);
403        let response = ApiVersionsResponseV0::decode_body(&mut decoder).unwrap();
404
405        assert_eq!(response.error_code, 0);
406        assert_eq!(response.api_keys.len(), 2);
407        assert_eq!(response.highest_supported_version(18, 3), Some(3));
408        assert_eq!(response.highest_supported_version(3, 12), Some(9));
409        assert_eq!(response.highest_supported_version(1, 1), None);
410        assert!(decoder.is_empty());
411    }
412
413    #[test]
414    fn encodes_api_versions_request_v3() {
415        let request = ApiVersionsRequestV3 {
416            correlation_id: 42,
417            client_id: Some("kafrust".to_owned()),
418            client_software_name: "kafrust".to_owned(),
419            client_software_version: "0.3.0".to_owned(),
420        };
421        assert_eq!(
422            request.encode().unwrap(),
423            [
424                0, 18, // api key
425                0, 3, // api version
426                0, 0, 0, 42, // correlation id
427                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', // nullable client id
428                0,    // request tagged fields
429                8, b'k', b'a', b'f', b'r', b'u', b's', b't', // software name
430                6, b'0', b'.', b'3', b'.', b'0', // software version
431                0,    // request body tagged fields
432            ]
433        );
434    }
435
436    #[test]
437    fn encodes_api_versions_request_v4_with_the_v4_header_version() {
438        let request = ApiVersionsRequestV4 {
439            correlation_id: 42,
440            client_id: Some("kafrust".to_owned()),
441            client_software_name: "kafrust".to_owned(),
442            client_software_version: "0.3.0".to_owned(),
443        };
444        assert_eq!(
445            request.encode().unwrap(),
446            [
447                0, 18, // api key
448                0, 4, // api version
449                0, 0, 0, 42, // correlation id
450                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', // nullable client id
451                0,    // request tagged fields
452                8, b'k', b'a', b'f', b'r', b'u', b's', b't', // software name
453                6, b'0', b'.', b'3', b'.', b'0', // software version
454                0,    // request body tagged fields
455            ]
456        );
457    }
458
459    #[test]
460    fn encodes_api_versions_request_v5_cluster_identity() {
461        let request = ApiVersionsRequestV5 {
462            correlation_id: 42,
463            client_id: None,
464            client_software_name: "kafrust".to_owned(),
465            client_software_version: "0.3.0".to_owned(),
466            cluster_id: Some("cluster-1".to_owned()),
467            node_id: 3,
468        };
469        assert_eq!(
470            request.encode().unwrap(),
471            [
472                0, 18, // api key
473                0, 5, // api version
474                0, 0, 0, 42, // correlation id
475                0xff, 0xff, // nullable client id
476                0,    // request header tagged fields
477                8, b'k', b'a', b'f', b'r', b'u', b's', b't', // software name
478                6, b'0', b'.', b'3', b'.', b'0', // software version
479                10, b'c', b'l', b'u', b's', b't', b'e', b'r', b'-', b'1', // cluster id
480                0, 0, 0, 3, // node id
481                0, // request body tagged fields
482            ]
483        );
484    }
485
486    #[test]
487    fn decodes_api_versions_response_v3() {
488        let bytes = [
489            0, 0, // error code
490            3, // compact api key count: two entries
491            0, 18, 0, 0, 0, 4, 0, // ApiVersions entry + tagged fields
492            0, 3, 0, 1, 0, 9, 0, // Metadata entry + tagged fields
493            0, 0, 0, 17, // throttle time
494            0,  // top-level tagged fields
495        ];
496        let mut decoder = Decoder::new(&bytes);
497        let response = ApiVersionsResponseV3::decode_body(&mut decoder).unwrap();
498
499        assert_eq!(response.error_code, 0);
500        assert_eq!(response.throttle_time_ms, 17);
501        assert_eq!(response.api_keys.len(), 2);
502        assert_eq!(response.highest_supported_version(18, 3), Some(3));
503        assert_eq!(response.highest_supported_version(3, 12), Some(9));
504        assert_eq!(response.highest_supported_version(1, 1), None);
505        assert!(response.tagged_fields.is_empty());
506        assert!(decoder.is_empty());
507    }
508
509    #[test]
510    fn decodes_api_versions_v4_feature_minimum_zero() {
511        let mut supported_features = Encoder::new();
512        supported_features
513            .write_compact_array(Some(&["metadata.version"]), |encoder, name| {
514                encoder.write_compact_string(name)?;
515                encoder.write_i16(0);
516                encoder.write_i16(1);
517                encoder.write_empty_tagged_fields();
518                Ok(())
519            })
520            .unwrap();
521
522        let supported_features = supported_features.into_bytes();
523        let mut body = Encoder::new();
524        body.write_i16(0);
525        body.write_unsigned_varint(1); // empty compact ApiKeys array
526        body.write_i32(0);
527        body.write_unsigned_varint(1); // one top-level tagged field
528        body.write_unsigned_varint(0); // supported features
529        body.write_unsigned_varint(u32::try_from(supported_features.len()).unwrap());
530        body.write_raw(&supported_features);
531
532        let bytes = body.into_bytes();
533        let mut decoder = Decoder::new(&bytes);
534        let response = ApiVersionsResponseV4::decode_body(&mut decoder).unwrap();
535
536        assert_eq!(response.supported_features.len(), 1);
537        assert_eq!(response.supported_features[0].min_version, 0);
538        assert_eq!(response.supported_features[0].max_version, 1);
539        assert!(decoder.is_empty());
540    }
541
542    #[test]
543    fn decodes_api_versions_feature_metadata_and_preserves_unknown_tags() {
544        let mut supported_features = Encoder::new();
545        supported_features
546            .write_compact_array(Some(&["group_coordinator"]), |encoder, name| {
547                encoder.write_compact_string(name)?;
548                encoder.write_i16(1);
549                encoder.write_i16(3);
550                encoder.write_empty_tagged_fields();
551                Ok(())
552            })
553            .unwrap();
554        let supported_features = supported_features.into_bytes();
555
556        let mut finalized_features = Encoder::new();
557        finalized_features
558            .write_compact_array(Some(&["metadata.version"]), |encoder, name| {
559                encoder.write_compact_string(name)?;
560                encoder.write_i16(4);
561                encoder.write_i16(1);
562                encoder.write_empty_tagged_fields();
563                Ok(())
564            })
565            .unwrap();
566        let finalized_features = finalized_features.into_bytes();
567
568        let mut body = Encoder::new();
569        body.write_i16(0);
570        body.write_unsigned_varint(1); // empty compact ApiKeys array
571        body.write_i32(0);
572        body.write_unsigned_varint(5); // four known tags and one unknown tag
573        body.write_unsigned_varint(0);
574        body.write_unsigned_varint(u32::try_from(supported_features.len()).unwrap());
575        body.write_raw(&supported_features);
576        body.write_unsigned_varint(1);
577        body.write_unsigned_varint(8);
578        body.write_i64(42);
579        body.write_unsigned_varint(2);
580        body.write_unsigned_varint(u32::try_from(finalized_features.len()).unwrap());
581        body.write_raw(&finalized_features);
582        body.write_unsigned_varint(3);
583        body.write_unsigned_varint(1);
584        body.write_bool(true);
585        body.write_unsigned_varint(99);
586        body.write_unsigned_varint(2);
587        body.write_raw(&[7, 8]);
588
589        let bytes = body.into_bytes();
590        let mut decoder = Decoder::new(&bytes);
591        let response = ApiVersionsResponseV3::decode_body(&mut decoder).unwrap();
592
593        assert_eq!(response.supported_features.len(), 1);
594        assert_eq!(response.supported_features[0].name, "group_coordinator");
595        assert_eq!(response.supported_features[0].min_version, 1);
596        assert_eq!(response.supported_features[0].max_version, 3);
597        assert_eq!(response.finalized_features_epoch, 42);
598        assert_eq!(response.finalized_features.len(), 1);
599        assert_eq!(response.finalized_features[0].name, "metadata.version");
600        assert_eq!(response.finalized_features[0].min_version_level, 1);
601        assert_eq!(response.finalized_features[0].max_version_level, 4);
602        assert!(response.zk_migration_ready);
603        assert_eq!(response.tagged_fields.len(), 1);
604        assert_eq!(response.tagged_fields[0].tag, 99);
605        assert_eq!(response.tagged_fields[0].data, vec![7, 8]);
606        assert!(decoder.is_empty());
607    }
608}