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 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#[derive(Debug, Clone, PartialEq, Eq)]
259pub struct ApiVersionsResponseV3 {
260 pub error_code: i16,
262 pub api_keys: Vec<ApiKeyVersion>,
264 pub throttle_time_ms: i32,
266 pub supported_features: Vec<SupportedFeature>,
268 pub finalized_features_epoch: i64,
270 pub finalized_features: Vec<FinalizedFeature>,
272 pub zk_migration_ready: bool,
274 pub tagged_fields: Vec<TaggedField>,
276}
277
278impl ApiVersionsResponseV3 {
279 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 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
344pub type ApiVersionsResponseV4 = ApiVersionsResponseV3;
346
347pub type ApiVersionsResponseV5 = ApiVersionsResponseV3;
349
350pub 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, 0, 0, 0, 0, 0, 42, 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, 0, 0, 0, 2, 0, 18, 0, 0, 0, 4, 0, 3, 0, 1, 0, 9, ];
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, 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, ]
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, 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, ]
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, 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, ]
483 );
484 }
485
486 #[test]
487 fn decodes_api_versions_response_v3() {
488 let bytes = [
489 0, 0, 3, 0, 18, 0, 0, 0, 4, 0, 0, 3, 0, 1, 0, 9, 0, 0, 0, 0, 17, 0, ];
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); body.write_i32(0);
527 body.write_unsigned_varint(1); body.write_unsigned_varint(0); 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); body.write_i32(0);
572 body.write_unsigned_varint(5); 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}