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#[derive(Debug, Clone, PartialEq, Eq)]
16pub struct SupportedFeature {
17 pub name: String,
19 pub min_version: i16,
21 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#[derive(Debug, Clone, PartialEq, Eq)]
39pub struct FinalizedFeature {
40 pub name: String,
42 pub min_version_level: i16,
44 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
60pub trait ApiVersionsLookup {
62 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#[derive(Debug, Clone, PartialEq, Eq)]
104pub struct ApiVersionsRequestV3 {
105 pub correlation_id: i32,
107 pub client_id: Option<String>,
109 pub client_software_name: String,
111 pub client_software_version: String,
113}
114
115impl ApiVersionsRequestV3 {
116 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#[derive(Debug, Clone, PartialEq, Eq)]
135pub struct ApiVersionsRequestV4 {
136 pub correlation_id: i32,
138 pub client_id: Option<String>,
140 pub client_software_name: String,
142 pub client_software_version: String,
144}
145
146impl ApiVersionsRequestV4 {
147 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#[derive(Debug, Clone, PartialEq, Eq)]
163pub struct ApiVersionsRequestV5 {
164 pub correlation_id: i32,
166 pub client_id: Option<String>,
168 pub client_software_name: String,
170 pub client_software_version: String,
172 pub cluster_id: Option<String>,
174 pub node_id: i32,
176}
177
178impl ApiVersionsRequestV5 {
179 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#[derive(Debug, Clone, PartialEq, Eq)]
257pub struct ApiVersionsResponseV3 {
258 pub error_code: i16,
260 pub api_keys: Vec<ApiKeyVersion>,
262 pub throttle_time_ms: i32,
264 pub supported_features: Vec<SupportedFeature>,
266 pub finalized_features_epoch: i64,
268 pub finalized_features: Vec<FinalizedFeature>,
270 pub zk_migration_ready: bool,
272 pub tagged_fields: Vec<TaggedField>,
274}
275
276impl ApiVersionsResponseV3 {
277 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 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
340pub type ApiVersionsResponseV4 = ApiVersionsResponseV3;
342
343pub type ApiVersionsResponseV5 = ApiVersionsResponseV3;
345
346pub 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, 0, 0, 0, 0, 0, 42, 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, 0, 0, 0, 2, 0, 18, 0, 0, 0, 4, 0, 3, 0, 1, 0, 9, ];
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, 0, 3, 0, 0, 0, 42, 0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', 0, 8, b'k', b'a', b'f', b'r', b'u', b's', b't', 6, b'0', b'.', b'3', b'.', b'0', 0, ]
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, 0, 4, 0, 0, 0, 42, 0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', 0, 8, b'k', b'a', b'f', b'r', b'u', b's', b't', 6, b'0', b'.', b'3', b'.', b'0', 0, ]
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, 0, 5, 0, 0, 0, 42, 0xff, 0xff, 0, 8, b'k', b'a', b'f', b'r', b'u', b's', b't', 6, b'0', b'.', b'3', b'.', b'0', 10, b'c', b'l', b'u', b's', b't', b'e', b'r', b'-', b'1', 0, 0, 0, 3, 0, ]
479 );
480 }
481
482 #[test]
483 fn decodes_api_versions_response_v3() {
484 let bytes = [
485 0, 0, 3, 0, 18, 0, 0, 0, 4, 0, 0, 3, 0, 1, 0, 9, 0, 0, 0, 0, 17, 0, ];
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); body.write_i32(0);
523 body.write_unsigned_varint(1); body.write_unsigned_varint(0); 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); body.write_i32(0);
568 body.write_unsigned_varint(5); 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}