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        Ok(Self {
233            error_code,
234            api_keys,
235        })
236    }
237
238    pub fn highest_supported_version(&self, api_key: i16, max_supported: i16) -> Option<i16> {
239        self.api_keys
240            .iter()
241            .find(|version| version.api_key == api_key)
242            .and_then(|version| {
243                let selected = version.max_version.min(max_supported);
244                (selected >= version.min_version).then_some(selected)
245            })
246    }
247}
248
249impl ApiVersionsLookup for ApiVersionsResponseV0 {
250    fn highest_supported_version(&self, api_key: i16, max_supported: i16) -> Option<i16> {
251        highest_supported_version(&self.api_keys, api_key, max_supported)
252    }
253}
254
255/// ApiVersions v3 response with flexible encoding and forward-compatible tags.
256#[derive(Debug, Clone, PartialEq, Eq)]
257pub struct ApiVersionsResponseV3 {
258    /// Top-level Kafka error code.
259    pub error_code: i16,
260    /// API version ranges advertised by the broker.
261    pub api_keys: Vec<ApiKeyVersion>,
262    /// Broker throttle duration in milliseconds.
263    pub throttle_time_ms: i32,
264    /// Feature ranges supported by the broker.
265    pub supported_features: Vec<SupportedFeature>,
266    /// Monotonically increasing finalized-feature metadata epoch, or `-1` if unknown.
267    pub finalized_features_epoch: i64,
268    /// Cluster-wide finalized feature ranges when the epoch is known.
269    pub finalized_features: Vec<FinalizedFeature>,
270    /// Whether the broker reports that ZooKeeper migration is ready.
271    pub zk_migration_ready: bool,
272    /// Unknown or future top-level tagged fields preserved for inspection.
273    pub tagged_fields: Vec<TaggedField>,
274}
275
276impl ApiVersionsResponseV3 {
277    /// Decodes an ApiVersions v3 response body after its response header.
278    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
279        let error_code = decoder.read_i16()?;
280        let api_keys = decoder
281            .read_compact_array("api versions", ApiKeyVersion::decode_flexible)?
282            .unwrap_or_default();
283        let throttle_time_ms = decoder.read_i32()?;
284        let tagged_fields = decoder.read_tagged_fields()?;
285        let limits = decoder.limits();
286        let mut supported_features = Vec::new();
287        let mut finalized_features_epoch = -1;
288        let mut finalized_features = Vec::new();
289        let mut zk_migration_ready = false;
290        let mut unknown_tagged_fields = Vec::new();
291        for field in tagged_fields {
292            match field.tag {
293                0 => {
294                    let mut field_decoder = Decoder::with_limits(&field.data, limits);
295                    supported_features = field_decoder
296                        .read_compact_array("supported features", SupportedFeature::decode)?
297                        .unwrap_or_default();
298                }
299                1 => {
300                    let mut field_decoder = Decoder::with_limits(&field.data, limits);
301                    finalized_features_epoch = field_decoder.read_i64()?;
302                }
303                2 => {
304                    let mut field_decoder = Decoder::with_limits(&field.data, limits);
305                    finalized_features = field_decoder
306                        .read_compact_array("finalized features", FinalizedFeature::decode)?
307                        .unwrap_or_default();
308                }
309                3 => {
310                    let mut field_decoder = Decoder::with_limits(&field.data, limits);
311                    zk_migration_ready = field_decoder.read_bool()?;
312                }
313                _ => unknown_tagged_fields.push(field),
314            }
315        }
316        Ok(Self {
317            error_code,
318            api_keys,
319            throttle_time_ms,
320            supported_features,
321            finalized_features_epoch,
322            finalized_features,
323            zk_migration_ready,
324            tagged_fields: unknown_tagged_fields,
325        })
326    }
327
328    /// Returns the highest broker-supported version not exceeding the client limit.
329    pub fn highest_supported_version(&self, api_key: i16, max_supported: i16) -> Option<i16> {
330        highest_supported_version(&self.api_keys, api_key, max_supported)
331    }
332}
333
334impl ApiVersionsLookup for ApiVersionsResponseV3 {
335    fn highest_supported_version(&self, api_key: i16, max_supported: i16) -> Option<i16> {
336        highest_supported_version(&self.api_keys, api_key, max_supported)
337    }
338}
339
340/// ApiVersions v4 response. The body schema is identical to v3.
341pub type ApiVersionsResponseV4 = ApiVersionsResponseV3;
342
343/// ApiVersions v5 response. The body schema is identical to v4.
344pub type ApiVersionsResponseV5 = ApiVersionsResponseV3;
345
346/// Kafka's `UNSUPPORTED_VERSION` protocol error code.
347pub const UNSUPPORTED_VERSION_ERROR_CODE: i16 = 35;
348
349fn highest_supported_version(
350    api_keys: &[ApiKeyVersion],
351    api_key: i16,
352    max_supported: i16,
353) -> Option<i16> {
354    api_keys
355        .iter()
356        .find(|version| version.api_key == api_key)
357        .and_then(|version| {
358            let selected = version.max_version.min(max_supported);
359            (selected >= version.min_version).then_some(selected)
360        })
361}
362
363#[cfg(test)]
364#[allow(clippy::unwrap_used)]
365mod tests {
366    use super::{
367        ApiVersionsRequestV0, ApiVersionsRequestV3, ApiVersionsRequestV4, ApiVersionsRequestV5,
368        ApiVersionsResponseV0, ApiVersionsResponseV3, ApiVersionsResponseV4, API_KEY,
369    };
370    use crate::codec::{Decoder, Encoder};
371
372    #[test]
373    fn encodes_api_versions_request_v0() {
374        let request = ApiVersionsRequestV0 {
375            correlation_id: 42,
376            client_id: Some("kafrust".to_owned()),
377        };
378        assert_eq!(
379            request.encode().unwrap(),
380            [
381                0, 18, // api key
382                0, 0, // api version
383                0, 0, 0, 42, // correlation id
384                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't',
385            ]
386        );
387        assert_eq!(API_KEY, 18);
388    }
389
390    #[test]
391    fn decodes_api_versions_response_v0() {
392        let bytes = [
393            0, 0, // error code
394            0, 0, 0, 2, // api key count
395            0, 18, 0, 0, 0, 4, // ApiVersions min/max
396            0, 3, 0, 1, 0, 9, // Metadata min/max
397        ];
398        let mut decoder = Decoder::new(&bytes);
399        let response = ApiVersionsResponseV0::decode_body(&mut decoder).unwrap();
400
401        assert_eq!(response.error_code, 0);
402        assert_eq!(response.api_keys.len(), 2);
403        assert_eq!(response.highest_supported_version(18, 3), Some(3));
404        assert_eq!(response.highest_supported_version(3, 12), Some(9));
405        assert_eq!(response.highest_supported_version(1, 1), None);
406        assert!(decoder.is_empty());
407    }
408
409    #[test]
410    fn encodes_api_versions_request_v3() {
411        let request = ApiVersionsRequestV3 {
412            correlation_id: 42,
413            client_id: Some("kafrust".to_owned()),
414            client_software_name: "kafrust".to_owned(),
415            client_software_version: "0.3.0".to_owned(),
416        };
417        assert_eq!(
418            request.encode().unwrap(),
419            [
420                0, 18, // api key
421                0, 3, // api version
422                0, 0, 0, 42, // correlation id
423                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', // nullable client id
424                0,    // request tagged fields
425                8, b'k', b'a', b'f', b'r', b'u', b's', b't', // software name
426                6, b'0', b'.', b'3', b'.', b'0', // software version
427                0,    // request body tagged fields
428            ]
429        );
430    }
431
432    #[test]
433    fn encodes_api_versions_request_v4_with_the_v4_header_version() {
434        let request = ApiVersionsRequestV4 {
435            correlation_id: 42,
436            client_id: Some("kafrust".to_owned()),
437            client_software_name: "kafrust".to_owned(),
438            client_software_version: "0.3.0".to_owned(),
439        };
440        assert_eq!(
441            request.encode().unwrap(),
442            [
443                0, 18, // api key
444                0, 4, // api version
445                0, 0, 0, 42, // correlation id
446                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', // nullable client id
447                0,    // request tagged fields
448                8, b'k', b'a', b'f', b'r', b'u', b's', b't', // software name
449                6, b'0', b'.', b'3', b'.', b'0', // software version
450                0,    // request body tagged fields
451            ]
452        );
453    }
454
455    #[test]
456    fn encodes_api_versions_request_v5_cluster_identity() {
457        let request = ApiVersionsRequestV5 {
458            correlation_id: 42,
459            client_id: None,
460            client_software_name: "kafrust".to_owned(),
461            client_software_version: "0.3.0".to_owned(),
462            cluster_id: Some("cluster-1".to_owned()),
463            node_id: 3,
464        };
465        assert_eq!(
466            request.encode().unwrap(),
467            [
468                0, 18, // api key
469                0, 5, // api version
470                0, 0, 0, 42, // correlation id
471                0xff, 0xff, // nullable client id
472                0,    // request header tagged fields
473                8, b'k', b'a', b'f', b'r', b'u', b's', b't', // software name
474                6, b'0', b'.', b'3', b'.', b'0', // software version
475                10, b'c', b'l', b'u', b's', b't', b'e', b'r', b'-', b'1', // cluster id
476                0, 0, 0, 3, // node id
477                0, // request body tagged fields
478            ]
479        );
480    }
481
482    #[test]
483    fn decodes_api_versions_response_v3() {
484        let bytes = [
485            0, 0, // error code
486            3, // compact api key count: two entries
487            0, 18, 0, 0, 0, 4, 0, // ApiVersions entry + tagged fields
488            0, 3, 0, 1, 0, 9, 0, // Metadata entry + tagged fields
489            0, 0, 0, 17, // throttle time
490            0,  // top-level tagged fields
491        ];
492        let mut decoder = Decoder::new(&bytes);
493        let response = ApiVersionsResponseV3::decode_body(&mut decoder).unwrap();
494
495        assert_eq!(response.error_code, 0);
496        assert_eq!(response.throttle_time_ms, 17);
497        assert_eq!(response.api_keys.len(), 2);
498        assert_eq!(response.highest_supported_version(18, 3), Some(3));
499        assert_eq!(response.highest_supported_version(3, 12), Some(9));
500        assert_eq!(response.highest_supported_version(1, 1), None);
501        assert!(response.tagged_fields.is_empty());
502        assert!(decoder.is_empty());
503    }
504
505    #[test]
506    fn decodes_api_versions_v4_feature_minimum_zero() {
507        let mut supported_features = Encoder::new();
508        supported_features
509            .write_compact_array(Some(&["metadata.version"]), |encoder, name| {
510                encoder.write_compact_string(name)?;
511                encoder.write_i16(0);
512                encoder.write_i16(1);
513                encoder.write_empty_tagged_fields();
514                Ok(())
515            })
516            .unwrap();
517
518        let supported_features = supported_features.into_bytes();
519        let mut body = Encoder::new();
520        body.write_i16(0);
521        body.write_unsigned_varint(1); // empty compact ApiKeys array
522        body.write_i32(0);
523        body.write_unsigned_varint(1); // one top-level tagged field
524        body.write_unsigned_varint(0); // supported features
525        body.write_unsigned_varint(u32::try_from(supported_features.len()).unwrap());
526        body.write_raw(&supported_features);
527
528        let bytes = body.into_bytes();
529        let mut decoder = Decoder::new(&bytes);
530        let response = ApiVersionsResponseV4::decode_body(&mut decoder).unwrap();
531
532        assert_eq!(response.supported_features.len(), 1);
533        assert_eq!(response.supported_features[0].min_version, 0);
534        assert_eq!(response.supported_features[0].max_version, 1);
535        assert!(decoder.is_empty());
536    }
537
538    #[test]
539    fn decodes_api_versions_feature_metadata_and_preserves_unknown_tags() {
540        let mut supported_features = Encoder::new();
541        supported_features
542            .write_compact_array(Some(&["group_coordinator"]), |encoder, name| {
543                encoder.write_compact_string(name)?;
544                encoder.write_i16(1);
545                encoder.write_i16(3);
546                encoder.write_empty_tagged_fields();
547                Ok(())
548            })
549            .unwrap();
550        let supported_features = supported_features.into_bytes();
551
552        let mut finalized_features = Encoder::new();
553        finalized_features
554            .write_compact_array(Some(&["metadata.version"]), |encoder, name| {
555                encoder.write_compact_string(name)?;
556                encoder.write_i16(4);
557                encoder.write_i16(1);
558                encoder.write_empty_tagged_fields();
559                Ok(())
560            })
561            .unwrap();
562        let finalized_features = finalized_features.into_bytes();
563
564        let mut body = Encoder::new();
565        body.write_i16(0);
566        body.write_unsigned_varint(1); // empty compact ApiKeys array
567        body.write_i32(0);
568        body.write_unsigned_varint(5); // four known tags and one unknown tag
569        body.write_unsigned_varint(0);
570        body.write_unsigned_varint(u32::try_from(supported_features.len()).unwrap());
571        body.write_raw(&supported_features);
572        body.write_unsigned_varint(1);
573        body.write_unsigned_varint(8);
574        body.write_i64(42);
575        body.write_unsigned_varint(2);
576        body.write_unsigned_varint(u32::try_from(finalized_features.len()).unwrap());
577        body.write_raw(&finalized_features);
578        body.write_unsigned_varint(3);
579        body.write_unsigned_varint(1);
580        body.write_bool(true);
581        body.write_unsigned_varint(99);
582        body.write_unsigned_varint(2);
583        body.write_raw(&[7, 8]);
584
585        let bytes = body.into_bytes();
586        let mut decoder = Decoder::new(&bytes);
587        let response = ApiVersionsResponseV3::decode_body(&mut decoder).unwrap();
588
589        assert_eq!(response.supported_features.len(), 1);
590        assert_eq!(response.supported_features[0].name, "group_coordinator");
591        assert_eq!(response.supported_features[0].min_version, 1);
592        assert_eq!(response.supported_features[0].max_version, 3);
593        assert_eq!(response.finalized_features_epoch, 42);
594        assert_eq!(response.finalized_features.len(), 1);
595        assert_eq!(response.finalized_features[0].name, "metadata.version");
596        assert_eq!(response.finalized_features[0].min_version_level, 1);
597        assert_eq!(response.finalized_features[0].max_version_level, 4);
598        assert!(response.zk_migration_ready);
599        assert_eq!(response.tagged_fields.len(), 1);
600        assert_eq!(response.tagged_fields[0].tag, 99);
601        assert_eq!(response.tagged_fields[0].data, vec![7, 8]);
602        assert!(decoder.is_empty());
603    }
604}