1use std::collections::{HashMap, HashSet};
4use std::fmt;
5use std::time::Duration;
6
7use bytes::{Buf, BufMut, Bytes, BytesMut};
8
9use super::api_keys::{pick_version, API_VERSIONS};
10use super::buf;
11use super::records::{self, RecordBatch};
12use crate::error::{Error, Result};
13use crate::net::BrokerConn;
14
15#[derive(Debug, Clone, PartialEq, Eq)]
17pub struct ApiVersion {
18 pub api_key: i16,
20 pub min_version: i16,
22 pub max_version: i16,
24}
25
26impl ApiVersion {
27 #[must_use]
29 pub fn api_key(&self) -> i16 {
30 self.api_key
31 }
32
33 #[must_use]
35 pub fn min_version(&self) -> i16 {
36 self.min_version
37 }
38
39 #[must_use]
41 pub fn max_version(&self) -> i16 {
42 self.max_version
43 }
44}
45
46#[derive(Debug, Clone, PartialEq, Eq)]
48pub struct SupportedFeatureKey {
49 pub name: String,
51 pub min_version: i16,
53 pub max_version: i16,
55}
56
57impl SupportedFeatureKey {
58 #[must_use]
60 pub fn name(&self) -> &str {
61 self.name.as_str()
62 }
63
64 #[must_use]
66 pub fn min_version(&self) -> i16 {
67 self.min_version
68 }
69
70 #[must_use]
72 pub fn max_version(&self) -> i16 {
73 self.max_version
74 }
75}
76
77#[derive(Debug, Clone, PartialEq, Eq)]
81pub struct FinalizedFeatureKey {
82 pub name: String,
84 pub max_version_level: i16,
86 pub min_version_level: i16,
88}
89
90impl FinalizedFeatureKey {
91 #[must_use]
93 pub fn name(&self) -> &str {
94 self.name.as_str()
95 }
96
97 #[must_use]
99 pub fn max_version_level(&self) -> i16 {
100 self.max_version_level
101 }
102
103 #[must_use]
105 pub fn min_version_level(&self) -> i16 {
106 self.min_version_level
107 }
108}
109
110#[derive(Debug, Clone, PartialEq, Eq, Default)]
112pub struct ApiVersionsResponse {
113 pub error_code: i16,
115 pub api_keys: Vec<ApiVersion>,
117 pub throttle_time_ms: i32,
119 pub supported_features: Vec<SupportedFeatureKey>,
121 pub finalized_features_epoch: Option<i64>,
123 pub finalized_features: Vec<FinalizedFeatureKey>,
125 pub zk_migration_ready: bool,
127}
128
129impl ApiVersionsResponse {
130 pub const UNKNOWN_FINALIZED_FEATURES_EPOCH: i64 = -1;
132
133 #[must_use]
135 pub fn api_version(&self, api_key: i16) -> Option<&ApiVersion> {
136 self.api_keys.iter().find(|k| k.api_key == api_key)
137 }
138
139 #[must_use]
141 pub const fn should_client_throttle(version: i16) -> bool {
142 version >= 2
143 }
144
145 #[must_use]
153 pub fn error_counts(&self) -> HashMap<i16, i32> {
154 HashMap::from([(self.error_code, 1)])
155 }
156
157 #[must_use]
159 pub fn zk_migration_ready(&self) -> bool {
160 self.zk_migration_ready
161 }
162
163 pub fn intersect(
169 this_version: Option<&ApiVersion>,
170 other: Option<&ApiVersion>,
171 ) -> Result<Option<ApiVersion>> {
172 let (Some(this_version), Some(other)) = (this_version, other) else {
173 return Ok(None);
174 };
175 if this_version.api_key() != other.api_key() {
176 return Err(Error::protocol(format!(
177 "thisVersion.apiKey: {} must be equal to other.apiKey: {}",
178 this_version.api_key(),
179 other.api_key()
180 )));
181 }
182 let min_version = this_version.min_version().max(other.min_version());
183 let max_version = this_version.max_version().min(other.max_version());
184 if min_version > max_version {
185 Ok(None)
186 } else {
187 Ok(Some(ApiVersion {
188 api_key: this_version.api_key(),
189 min_version,
190 max_version,
191 }))
192 }
193 }
194
195 #[must_use]
203 pub fn create_finalized_feature_keys<'a, I>(finalized_features: I) -> Vec<FinalizedFeatureKey>
204 where
205 I: IntoIterator<Item = (&'a str, i16)>,
206 {
207 let mut order: Vec<String> = Vec::new();
208 let mut levels: HashMap<String, i16> = HashMap::new();
209 for (name, level) in finalized_features {
210 match levels.entry(name.to_string()) {
211 std::collections::hash_map::Entry::Vacant(slot) => {
212 order.push(slot.key().clone());
213 let _inserted = slot.insert(level);
214 }
215 std::collections::hash_map::Entry::Occupied(mut slot) => {
216 let _prev = slot.insert(level);
217 }
218 }
219 }
220 order
221 .into_iter()
222 .filter_map(|name| {
223 let level = levels.remove(&name)?;
224 (level != 0).then_some(FinalizedFeatureKey {
225 name,
226 max_version_level: level,
227 min_version_level: level,
228 })
229 })
230 .collect()
231 }
232
233 #[must_use]
242 pub fn maybe_filter_supported_feature_keys(
243 supported_features: &[SupportedFeatureKey],
244 alter_feature_level_0: bool,
245 ) -> Vec<SupportedFeatureKey> {
246 supported_features
247 .iter()
248 .filter(|feature| !(alter_feature_level_0 && feature.min_version == 0))
249 .cloned()
250 .collect()
251 }
252}
253
254pub struct ApiVersionsRequest;
256
257impl ApiVersionsRequest {
258 #[must_use]
265 pub fn is_valid(version: i16, software_name: &str, software_version: &str) -> bool {
266 version < 3
267 || (client_software_name_or_version_ok(software_name)
268 && client_software_name_or_version_ok(software_version))
269 }
270
271 #[must_use]
282 pub fn error_response(error_code: i16) -> ApiVersionsResponse {
283 ApiVersionsResponse {
284 error_code,
285 api_keys: if error_code == crate::error::UNSUPPORTED_VERSION {
286 vec![ApiVersion {
287 api_key: API_VERSIONS,
288 min_version: 0,
289 max_version: 4,
290 }]
291 } else {
292 Vec::new()
293 },
294 ..Default::default()
295 }
296 }
297}
298
299fn client_software_name_or_version_ok(s: &str) -> bool {
301 let mut chars = s.chars();
302 let Some(first) = chars.next() else {
303 return false;
304 };
305 if !first.is_ascii_alphanumeric() {
306 return false;
307 }
308 let Some(last) = chars.next_back() else {
309 return true;
310 };
311 last.is_ascii_alphanumeric() && chars.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '.')
312}
313
314fn api_versions_flexible(version: i16) -> Result<bool> {
321 match version {
322 0..=2 => Ok(false),
323 3..=4 => Ok(true),
324 other => Err(Error::protocol(format!(
325 "ApiVersions version {other} is not implemented"
326 ))),
327 }
328}
329
330pub fn encode_api_versions_request(
337 buf: &mut BytesMut,
338 version: i16,
339 software_name: &str,
340 software_version: &str,
341) -> crate::error::Result<()> {
342 let flexible = api_versions_flexible(version)?;
343 if version >= 3 {
344 buf::put_string(buf, flexible, Some(software_name))?;
345 buf::put_string(buf, flexible, Some(software_version))?;
346 buf::put_empty_tagged_fields(buf);
347 }
348 Ok(())
349}
350
351pub fn decode_api_versions_request<B: Buf>(buf: &mut B, version: i16) -> Result<(String, String)> {
354 let flexible = api_versions_flexible(version)?;
355 if version >= 3 {
356 let name = buf::get_string(buf, flexible)?.unwrap_or_default();
357 let software_version = buf::get_string(buf, flexible)?.unwrap_or_default();
358 buf::skip_tagged_fields(buf)?;
359 return Ok((name, software_version));
360 }
361 Ok((String::new(), String::new()))
362}
363
364pub fn decode_api_versions_response<B: Buf>(
366 buf: &mut B,
367 version: i16,
368) -> Result<ApiVersionsResponse> {
369 let flexible = api_versions_flexible(version)?;
370 let error_code = buf::get_i16(buf)?;
371 let count = buf::get_array_len(buf, flexible)?.unwrap_or(0);
372 let mut api_keys = Vec::with_capacity(count);
373 for _ in 0..count {
374 let api_key = buf::get_i16(buf)?;
375 let min_version = buf::get_i16(buf)?;
376 let max_version = buf::get_i16(buf)?;
377 if flexible {
378 buf::skip_tagged_fields(buf)?;
379 }
380 api_keys.push(ApiVersion {
381 api_key,
382 min_version,
383 max_version,
384 });
385 }
386 let throttle_time_ms = if version >= 1 { buf::get_i32(buf)? } else { 0 };
387 let tagged = if flexible {
388 decode_api_versions_tagged_fields(buf)?
389 } else {
390 ApiVersionsTaggedFields {
391 supported: Vec::new(),
392 epoch: None,
393 finalized: Vec::new(),
394 zk_migration_ready: false,
395 }
396 };
397 Ok(ApiVersionsResponse {
398 error_code,
399 api_keys,
400 throttle_time_ms,
401 supported_features: tagged.supported,
402 finalized_features_epoch: tagged.epoch,
403 finalized_features: tagged.finalized,
404 zk_migration_ready: tagged.zk_migration_ready,
405 })
406}
407
408pub fn decode_api_versions_handshake(bytes: &[u8], version: i16) -> Result<ApiVersionsResponse> {
412 let mut cur = bytes;
413 if let Ok(resp) = decode_api_versions_response(&mut cur, version) {
414 if !cur.has_remaining() {
415 return Ok(resp);
416 }
417 }
418 if version != 0 {
419 let mut cur = bytes;
420 if let Ok(resp) = decode_api_versions_response(&mut cur, 0) {
421 if !cur.has_remaining() {
422 return Ok(resp);
423 }
424 }
425 }
426 let mut cur = bytes;
427 decode_api_versions_response(&mut cur, version)
428}
429
430pub async fn negotiate_api_versions(
437 conn: &mut BrokerConn,
438 request_timeout: Duration,
439) -> Result<ApiVersionsResponse> {
440 const SENT: i16 = 4;
441 let body = conn
442 .roundtrip(
443 API_VERSIONS,
444 SENT,
445 |buf| encode_api_versions_request(buf, SENT, "partitionline", "0.1.0"),
446 request_timeout,
447 )
448 .await?;
449 let resp = decode_api_versions_handshake(body.as_ref(), SENT)?;
450 if resp.error_code == 0 {
451 return Ok(resp);
452 }
453 if resp.error_code != crate::error::UNSUPPORTED_VERSION {
454 return Err(Error::broker(resp.error_code, "ApiVersions"));
455 }
456 let retry = resp
457 .api_version(API_VERSIONS)
458 .and_then(|v| pick_version(v.min_version, v.max_version, 0, SENT))
459 .unwrap_or(0);
460 if retry == SENT {
461 return Err(Error::broker(resp.error_code, "ApiVersions"));
462 }
463 let body = conn
464 .roundtrip(
465 API_VERSIONS,
466 retry,
467 |buf| encode_api_versions_request(buf, retry, "partitionline", "0.1.0"),
468 request_timeout,
469 )
470 .await?;
471 let resp = decode_api_versions_handshake(body.as_ref(), retry)?;
472 if resp.error_code != 0 {
473 return Err(Error::broker(resp.error_code, "ApiVersions"));
474 }
475 Ok(resp)
476}
477
478pub fn encode_api_versions_response(
482 buf: &mut BytesMut,
483 version: i16,
484 resp: &ApiVersionsResponse,
485) -> crate::error::Result<()> {
486 let flexible = api_versions_flexible(version)?;
487 buf.put_i16(resp.error_code);
488 buf::put_array_len(buf, flexible, Some(resp.api_keys.len()))?;
489 for api in &resp.api_keys {
490 buf.put_i16(api.api_key);
491 buf.put_i16(api.min_version);
492 buf.put_i16(api.max_version);
493 if flexible {
494 buf::put_empty_tagged_fields(buf);
495 }
496 }
497 if version >= 1 {
498 buf.put_i32(resp.throttle_time_ms);
499 }
500 if flexible {
501 encode_api_versions_tagged_fields(buf, version, resp)?;
502 }
503 Ok(())
504}
505
506struct ApiVersionsTaggedFields {
507 supported: Vec<SupportedFeatureKey>,
508 epoch: Option<i64>,
509 finalized: Vec<FinalizedFeatureKey>,
510 zk_migration_ready: bool,
511}
512
513fn leftover_empty<B: Buf>(buf: &B, what: &'static str) -> Result<()> {
514 if buf.has_remaining() {
515 Err(Error::protocol(format!("{what} leftover")))
516 } else {
517 Ok(())
518 }
519}
520
521fn encode_supported_feature_keys(features: &[SupportedFeatureKey]) -> Result<Bytes> {
522 let mut buf = BytesMut::new();
523 buf::put_array_len(&mut buf, true, Some(features.len()))?;
524 for f in features {
525 buf::put_compact_string(&mut buf, Some(&f.name))?;
526 buf.put_i16(f.min_version);
527 buf.put_i16(f.max_version);
528 buf::put_empty_tagged_fields(&mut buf);
529 }
530 Ok(buf.freeze())
531}
532
533fn decode_supported_feature_keys<B: Buf>(buf: &mut B) -> Result<Vec<SupportedFeatureKey>> {
534 let n = buf::get_array_len(buf, true)?.unwrap_or(0);
535 let mut out = Vec::with_capacity(n);
536 for _ in 0..n {
537 let name = buf::get_compact_string(buf)?.unwrap_or_default();
538 let min_version = buf::get_i16(buf)?;
539 let max_version = buf::get_i16(buf)?;
540 buf::skip_tagged_fields(buf)?;
541 out.push(SupportedFeatureKey {
542 name,
543 min_version,
544 max_version,
545 });
546 }
547 Ok(out)
548}
549
550fn encode_finalized_feature_keys(features: &[FinalizedFeatureKey]) -> Result<Bytes> {
551 let mut buf = BytesMut::new();
552 buf::put_array_len(&mut buf, true, Some(features.len()))?;
553 for f in features {
554 buf::put_compact_string(&mut buf, Some(&f.name))?;
555 buf.put_i16(f.max_version_level);
556 buf.put_i16(f.min_version_level);
557 buf::put_empty_tagged_fields(&mut buf);
558 }
559 Ok(buf.freeze())
560}
561
562fn decode_finalized_feature_keys<B: Buf>(buf: &mut B) -> Result<Vec<FinalizedFeatureKey>> {
563 let n = buf::get_array_len(buf, true)?.unwrap_or(0);
564 let mut out = Vec::with_capacity(n);
565 for _ in 0..n {
566 let name = buf::get_compact_string(buf)?.unwrap_or_default();
567 let max_version_level = buf::get_i16(buf)?;
568 let min_version_level = buf::get_i16(buf)?;
569 buf::skip_tagged_fields(buf)?;
570 out.push(FinalizedFeatureKey {
571 name,
572 max_version_level,
573 min_version_level,
574 });
575 }
576 Ok(out)
577}
578
579fn encode_api_versions_tagged_fields(
580 buf: &mut BytesMut,
581 version: i16,
582 resp: &ApiVersionsResponse,
583) -> Result<()> {
584 let mut tags: Vec<(u32, Bytes)> = Vec::new();
585 let supported: Vec<SupportedFeatureKey> = resp
586 .supported_features
587 .iter()
588 .filter(|f| version >= 4 || f.min_version != 0)
589 .cloned()
590 .collect();
591 if !supported.is_empty() {
592 tags.push((0, encode_supported_feature_keys(&supported)?));
593 }
594 if let Some(epoch) = resp.finalized_features_epoch {
595 if epoch >= 0 {
596 let mut b = BytesMut::new();
597 b.put_i64(epoch);
598 tags.push((1, b.freeze()));
599 }
600 }
601 if !resp.finalized_features.is_empty() {
602 tags.push((2, encode_finalized_feature_keys(&resp.finalized_features)?));
603 }
604 if resp.zk_migration_ready {
605 tags.push((3, Bytes::from_static(&[1])));
606 }
607 buf::put_tagged_fields(buf, &tags)
608}
609
610fn decode_api_versions_tagged_fields<B: Buf>(buf: &mut B) -> Result<ApiVersionsTaggedFields> {
611 let tags = buf::get_tagged_fields(buf)?;
612 let mut supported = Vec::new();
613 let mut epoch = None;
614 let mut finalized = Vec::new();
615 let mut zk = false;
616 for (tag, value) in tags {
617 match tag {
618 0 => {
619 let mut cur = value.as_ref();
620 supported = decode_supported_feature_keys(&mut cur)?;
621 leftover_empty(&cur, "supported_features")?;
622 }
623 1 => {
624 let mut cur = value.as_ref();
625 let v = buf::get_i64(&mut cur)?;
626 leftover_empty(&cur, "finalized_features_epoch")?;
627 epoch = (v != ApiVersionsResponse::UNKNOWN_FINALIZED_FEATURES_EPOCH && v >= 0)
628 .then_some(v);
629 }
630 2 => {
631 let mut cur = value.as_ref();
632 finalized = decode_finalized_feature_keys(&mut cur)?;
633 leftover_empty(&cur, "finalized_features")?;
634 }
635 3 => {
636 let mut cur = value.as_ref();
637 zk = buf::get_bool(&mut cur)?;
638 leftover_empty(&cur, "zk_migration_ready")?;
639 }
640 _ => {}
641 }
642 }
643 Ok(ApiVersionsTaggedFields {
644 supported,
645 epoch,
646 finalized,
647 zk_migration_ready: zk,
648 })
649}
650
651pub(crate) fn format_java_node(
653 f: &mut fmt::Formatter<'_>,
654 host: &str,
655 port: i32,
656 id: i32,
657 rack: Option<&str>,
658 is_fenced: bool,
659) -> fmt::Result {
660 write!(
661 f,
662 "{}:{} (id: {} rack: {} isFenced: {})",
663 host,
664 port,
665 id,
666 rack.unwrap_or("null"),
667 is_fenced
668 )
669}
670
671#[derive(Debug, Clone, PartialEq, Eq)]
673pub struct Broker {
674 pub node_id: i32,
676 pub host: String,
678 pub port: i32,
680 pub rack: Option<String>,
682}
683
684impl Broker {
685 pub fn new(node_id: i32, host: impl Into<String>, port: i32, rack: Option<String>) -> Self {
687 Self {
688 node_id,
689 host: host.into(),
690 port,
691 rack,
692 }
693 }
694
695 #[must_use]
697 pub fn id(&self) -> i32 {
698 self.node_id
699 }
700
701 #[must_use]
703 pub fn host(&self) -> &str {
704 self.host.as_str()
705 }
706
707 #[must_use]
709 pub fn port(&self) -> i32 {
710 self.port
711 }
712
713 #[must_use]
715 pub fn rack(&self) -> Option<&str> {
716 self.rack.as_deref()
717 }
718
719 #[must_use]
721 pub fn has_rack(&self) -> bool {
722 self.rack.is_some()
723 }
724
725 #[must_use]
727 pub fn is_fenced(&self) -> bool {
728 false
729 }
730
731 #[must_use]
733 pub fn id_string(&self) -> String {
734 self.node_id.to_string()
735 }
736
737 #[must_use]
739 pub fn is_empty(&self) -> bool {
740 self.host.is_empty() || self.port < 0
741 }
742
743 #[must_use]
745 pub fn no_node() -> Self {
746 Self::new(-1, "", -1, None)
747 }
748}
749
750impl fmt::Display for Broker {
751 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
752 format_java_node(
753 f,
754 self.host.as_str(),
755 self.port,
756 self.node_id,
757 self.rack.as_deref(),
758 false,
759 )
760 }
761}
762
763#[derive(Debug, Clone, PartialEq, Eq)]
765pub struct NodeEndpoint {
766 pub node_id: i32,
768 pub host: String,
770 pub port: i32,
772 pub rack: Option<String>,
774}
775
776impl NodeEndpoint {
777 pub fn new(node_id: i32, host: impl Into<String>, port: i32, rack: Option<String>) -> Self {
779 Self {
780 node_id,
781 host: host.into(),
782 port,
783 rack,
784 }
785 }
786
787 #[must_use]
789 pub fn id(&self) -> i32 {
790 self.node_id
791 }
792
793 #[must_use]
795 pub fn host(&self) -> &str {
796 self.host.as_str()
797 }
798
799 #[must_use]
801 pub fn port(&self) -> i32 {
802 self.port
803 }
804
805 #[must_use]
807 pub fn rack(&self) -> Option<&str> {
808 self.rack.as_deref()
809 }
810
811 #[must_use]
813 pub fn has_rack(&self) -> bool {
814 self.rack.is_some()
815 }
816
817 #[must_use]
819 pub fn is_fenced(&self) -> bool {
820 false
821 }
822
823 #[must_use]
825 pub fn id_string(&self) -> String {
826 self.node_id.to_string()
827 }
828
829 #[must_use]
831 pub fn is_empty(&self) -> bool {
832 self.host.is_empty() || self.port < 0
833 }
834
835 #[must_use]
837 pub fn no_node() -> Self {
838 Self::new(-1, "", -1, None)
839 }
840}
841
842impl fmt::Display for NodeEndpoint {
843 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
844 format_java_node(
845 f,
846 self.host.as_str(),
847 self.port,
848 self.node_id,
849 self.rack.as_deref(),
850 false,
851 )
852 }
853}
854
855impl From<Broker> for NodeEndpoint {
856 fn from(b: Broker) -> Self {
857 Self {
858 node_id: b.node_id,
859 host: b.host,
860 port: b.port,
861 rack: b.rack,
862 }
863 }
864}
865
866impl From<NodeEndpoint> for Broker {
867 fn from(e: NodeEndpoint) -> Self {
868 Self {
869 node_id: e.node_id,
870 host: e.host,
871 port: e.port,
872 rack: e.rack,
873 }
874 }
875}
876
877pub(crate) fn put_compact_node_endpoints(
882 buf: &mut BytesMut,
883 endpoints: &[NodeEndpoint],
884) -> Result<()> {
885 buf::put_array_len(buf, true, Some(endpoints.len()))?;
886 for e in endpoints {
887 buf.put_i32(e.node_id);
888 buf::put_string(buf, true, Some(e.host.as_str()))?;
889 buf.put_i32(e.port);
890 buf::put_string(buf, true, e.rack.as_deref())?;
891 buf::put_empty_tagged_fields(buf);
892 }
893 Ok(())
894}
895
896pub(crate) fn get_compact_node_endpoints<B: Buf>(buf: &mut B) -> Result<Vec<NodeEndpoint>> {
898 let n = buf::get_array_len(buf, true)?.unwrap_or(0);
899 let mut out = Vec::with_capacity(n);
900 for _ in 0..n {
901 let node_id = buf::get_i32(buf)?;
902 let host = buf::get_string(buf, true)?.unwrap_or_default();
903 let port = buf::get_i32(buf)?;
904 let rack = buf::get_string(buf, true)?;
905 buf::skip_tagged_fields(buf)?;
906 out.push(NodeEndpoint {
907 node_id,
908 host,
909 port,
910 rack,
911 });
912 }
913 Ok(out)
914}
915
916pub(crate) fn encode_node_endpoints(endpoints: &[NodeEndpoint]) -> Result<Bytes> {
918 let mut inner = BytesMut::new();
919 put_compact_node_endpoints(&mut inner, endpoints)?;
920 Ok(inner.freeze())
921}
922
923pub(crate) fn decode_node_endpoints(value: &Bytes) -> Result<Vec<NodeEndpoint>> {
925 let mut cur = value.as_ref();
926 let out = get_compact_node_endpoints(&mut cur)?;
927 leftover_empty(&cur, "NodeEndpoints")?;
928 Ok(out)
929}
930
931pub(crate) fn encode_top_level_node_endpoints(
932 buf: &mut BytesMut,
933 include: bool,
934 endpoints: &[NodeEndpoint],
935) -> Result<()> {
936 if include && !endpoints.is_empty() {
937 buf::put_tagged_fields(buf, &[(0, encode_node_endpoints(endpoints)?)])
938 } else {
939 buf::put_empty_tagged_fields(buf);
940 Ok(())
941 }
942}
943
944pub(crate) fn decode_top_level_node_endpoints<B: Buf>(
945 buf: &mut B,
946 include: bool,
947) -> Result<Vec<NodeEndpoint>> {
948 let tags = buf::get_tagged_fields(buf)?;
949 let mut endpoints = Vec::new();
950 for (tag, value) in tags {
951 if include && tag == 0 {
952 endpoints = decode_node_endpoints(&value)?;
953 }
954 }
955 Ok(endpoints)
956}
957
958#[derive(Debug, Clone, PartialEq, Eq)]
967pub struct PartitionMetadata {
968 pub error_code: i16,
970 pub partition_index: i32,
972 pub leader_id: i32,
974 pub leader_epoch: i32,
976 pub replica_nodes: Vec<i32>,
978 pub isr_nodes: Vec<i32>,
980 pub offline_replicas: Vec<i32>,
982}
983
984impl PartitionMetadata {
985 #[must_use]
999 pub fn new(
1000 error_code: i16,
1001 partition_index: i32,
1002 leader_id: Option<i32>,
1003 leader_epoch: Option<i32>,
1004 replica_nodes: Vec<i32>,
1005 isr_nodes: Vec<i32>,
1006 offline_replicas: Vec<i32>,
1007 ) -> Self {
1008 Self {
1009 error_code,
1010 partition_index,
1011 leader_id: leader_id.unwrap_or(MetadataResponse::NO_LEADER_ID),
1012 leader_epoch: leader_epoch.unwrap_or(RecordBatch::NO_PARTITION_LEADER_EPOCH),
1013 replica_nodes,
1014 isr_nodes,
1015 offline_replicas,
1016 }
1017 }
1018
1019 #[must_use]
1024 pub fn without_leader_epoch(&self) -> Self {
1025 Self {
1026 leader_epoch: RecordBatch::NO_PARTITION_LEADER_EPOCH,
1027 ..self.clone()
1028 }
1029 }
1030}
1031
1032#[derive(Debug, Clone, PartialEq, Eq)]
1046pub struct TopicMetadata {
1047 pub error_code: i16,
1049 pub name: Option<String>,
1051 pub topic_id: [u8; 16],
1053 pub is_internal: bool,
1055 pub partitions: Vec<PartitionMetadata>,
1057 pub topic_authorized_operations: i32,
1060}
1061
1062impl TopicMetadata {
1063 #[must_use]
1070 pub fn error(error_code: i16, name: Option<&str>, topic_id: [u8; 16]) -> Self {
1071 Self {
1072 error_code,
1073 name: Some(name.map(str::to_owned).unwrap_or_default()),
1074 topic_id,
1075 is_internal: false,
1076 partitions: Vec::new(),
1077 topic_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
1078 }
1079 }
1080
1081 #[must_use]
1094 pub fn new(
1095 error_code: i16,
1096 topic: impl Into<String>,
1097 is_internal: bool,
1098 partitions: Vec<PartitionMetadata>,
1099 ) -> Self {
1100 Self::with_topic_id(
1101 error_code,
1102 topic,
1103 [0; 16],
1104 is_internal,
1105 partitions,
1106 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
1107 )
1108 }
1109
1110 #[must_use]
1121 pub fn with_topic_id(
1122 error_code: i16,
1123 topic: impl Into<String>,
1124 topic_id: [u8; 16],
1125 is_internal: bool,
1126 partitions: Vec<PartitionMetadata>,
1127 topic_authorized_operations: i32,
1128 ) -> Self {
1129 Self {
1130 error_code,
1131 name: Some(topic.into()),
1132 topic_id,
1133 is_internal,
1134 partitions,
1135 topic_authorized_operations,
1136 }
1137 }
1138}
1139
1140fn write_java_uuid_bytes(f: &mut fmt::Formatter<'_>, bytes: &[u8; 16]) -> fmt::Result {
1141 f.write_str(&base64::Engine::encode(
1142 &base64::engine::general_purpose::URL_SAFE_NO_PAD,
1143 bytes.as_slice(),
1144 ))
1145}
1146
1147fn write_java_optional_i32(f: &mut fmt::Formatter<'_>, value: Option<i32>) -> fmt::Result {
1148 match value {
1149 None => f.write_str("Optional.empty"),
1150 Some(n) => write!(f, "Optional[{n}]"),
1151 }
1152}
1153
1154fn write_java_int_csv(f: &mut fmt::Formatter<'_>, ids: &[i32]) -> fmt::Result {
1155 for (i, id) in ids.iter().enumerate() {
1156 if i > 0 {
1157 f.write_str(",")?;
1158 }
1159 write!(f, "{id}")?;
1160 }
1161 Ok(())
1162}
1163
1164fn write_java_partition_metadata(
1165 f: &mut fmt::Formatter<'_>,
1166 topic: Option<&str>,
1167 p: &PartitionMetadata,
1168) -> fmt::Result {
1169 f.write_str("PartitionMetadata(error=")?;
1170 f.write_str(crate::error::for_code(p.error_code))?;
1171 f.write_str(", partition=")?;
1172 match topic {
1173 Some(name) => write!(f, "{name}-{}", p.partition_index)?,
1174 None => write!(f, "null-{}", p.partition_index)?,
1175 }
1176 f.write_str(", leader=")?;
1177 write_java_optional_i32(f, (p.leader_id >= 0).then_some(p.leader_id))?;
1178 f.write_str(", leaderEpoch=")?;
1179 write_java_optional_i32(
1180 f,
1181 (p.leader_epoch != RecordBatch::NO_PARTITION_LEADER_EPOCH).then_some(p.leader_epoch),
1182 )?;
1183 f.write_str(", replicas=")?;
1184 write_java_int_csv(f, &p.replica_nodes)?;
1185 f.write_str(", isr=")?;
1186 write_java_int_csv(f, &p.isr_nodes)?;
1187 f.write_str(", offlineReplicas=")?;
1188 write_java_int_csv(f, &p.offline_replicas)?;
1189 f.write_str(")")
1190}
1191
1192impl fmt::Display for TopicMetadata {
1193 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1194 f.write_str("TopicMetadata{error=")?;
1195 f.write_str(crate::error::for_code(self.error_code))?;
1196 f.write_str(", topic='")?;
1197 match self.name.as_deref() {
1198 Some(name) => f.write_str(name)?,
1199 None => f.write_str("null")?,
1200 }
1201 f.write_str("', topicId='")?;
1202 write_java_uuid_bytes(f, &self.topic_id)?;
1203 write!(f, "', isInternal={}, partitionMetadata=[", self.is_internal)?;
1204 for (i, p) in self.partitions.iter().enumerate() {
1205 if i > 0 {
1206 f.write_str(", ")?;
1207 }
1208 write_java_partition_metadata(f, self.name.as_deref(), p)?;
1209 }
1210 write!(
1211 f,
1212 "], authorizedOperations={}}}",
1213 self.topic_authorized_operations
1214 )
1215 }
1216}
1217
1218#[derive(Debug, Clone, PartialEq, Eq)]
1220pub struct MetadataResponse {
1221 pub throttle_time_ms: i32,
1223 pub brokers: Vec<Broker>,
1225 pub cluster_id: Option<String>,
1227 pub controller_id: i32,
1229 pub topics: Vec<TopicMetadata>,
1231 pub cluster_authorized_operations: i32,
1234 pub error_code: i16,
1236}
1237
1238impl MetadataResponse {
1239 pub const NO_CONTROLLER_ID: i32 = -1;
1241 pub const NO_LEADER_ID: i32 = -1;
1243 pub const AUTHORIZED_OPERATIONS_OMITTED: i32 = i32::MIN;
1245
1246 #[must_use]
1253 pub const fn has_reliable_leader_epochs(version: i16) -> bool {
1254 version >= 9
1255 }
1256
1257 #[must_use]
1259 pub const fn should_client_throttle(version: i16) -> bool {
1260 version >= 6
1261 }
1262
1263 pub fn errors(&self) -> Result<HashMap<String, i16>> {
1270 let mut errors = HashMap::new();
1271 for topic in &self.topics {
1272 let Some(name) = topic.name.as_ref() else {
1273 return Err(Error::protocol(
1274 "Use errorsByTopicId() when managing topic using topic id",
1275 ));
1276 };
1277 if topic.error_code != 0 {
1278 let _prev = errors.insert(name.clone(), topic.error_code);
1279 }
1280 }
1281 Ok(errors)
1282 }
1283
1284 pub fn errors_by_topic_id(&self) -> Result<HashMap<[u8; 16], i16>> {
1291 let mut errors = HashMap::new();
1292 for topic in &self.topics {
1293 if topic.topic_id == [0u8; 16] {
1294 return Err(Error::protocol(
1295 "Use errors() when managing topic using topic name",
1296 ));
1297 }
1298 if topic.error_code != 0 {
1299 let _prev = errors.insert(topic.topic_id, topic.error_code);
1300 }
1301 }
1302 Ok(errors)
1303 }
1304
1305 #[must_use]
1310 pub fn topics_by_error(&self, error_code: i16) -> HashSet<String> {
1311 self.topics
1312 .iter()
1313 .filter(|topic| topic.error_code == error_code)
1314 .filter_map(|topic| topic.name.clone())
1315 .collect()
1316 }
1317
1318 #[must_use]
1323 pub fn error_counts(&self) -> HashMap<i16, i32> {
1324 let mut counts = HashMap::new();
1325 for topic in &self.topics {
1326 for partition in &topic.partitions {
1327 let count = counts.entry(partition.error_code).or_insert(0);
1328 *count += 1;
1329 }
1330 let count = counts.entry(topic.error_code).or_insert(0);
1331 *count += 1;
1332 }
1333 counts
1334 }
1335
1336 #[must_use]
1340 pub fn topic_authorized_operations(&self, topic_name: &str) -> Option<i32> {
1341 self.topics
1342 .iter()
1343 .find(|topic| topic.name.as_deref() == Some(topic_name))
1344 .map(|topic| topic.topic_authorized_operations)
1345 }
1346
1347 #[must_use]
1349 pub fn cluster_authorized_operations(&self) -> i32 {
1350 self.cluster_authorized_operations
1351 }
1352
1353 #[must_use]
1355 pub fn brokers_by_id(&self) -> HashMap<i32, Broker> {
1356 self.brokers
1357 .iter()
1358 .map(|broker| (broker.node_id, broker.clone()))
1359 .collect()
1360 }
1361
1362 #[must_use]
1371 pub fn controller(&self) -> Option<&Broker> {
1372 self.brokers
1373 .iter()
1374 .rfind(|broker| broker.node_id == self.controller_id)
1375 }
1376
1377 pub(crate) fn check(&self) -> Result<()> {
1379 if self.error_code == 0 {
1380 Ok(())
1381 } else {
1382 Err(Error::broker(self.error_code, "Metadata"))
1383 }
1384 }
1385}
1386
1387#[derive(Debug, Clone, PartialEq, Eq)]
1398pub struct MetadataRequestTopic {
1399 pub name: Option<String>,
1401 pub topic_id: [u8; 16],
1403}
1404
1405impl MetadataRequestTopic {
1406 #[must_use]
1408 pub fn by_name(name: impl Into<String>) -> Self {
1409 Self {
1410 name: Some(name.into()),
1411 topic_id: [0; 16],
1412 }
1413 }
1414
1415 #[must_use]
1419 pub fn by_id(topic_id: [u8; 16]) -> Self {
1420 Self {
1421 name: None,
1422 topic_id,
1423 }
1424 }
1425
1426 #[must_use]
1428 pub fn convert_from_names<I>(names: I) -> Vec<Self>
1429 where
1430 I: IntoIterator,
1431 I::Item: Into<String>,
1432 {
1433 names.into_iter().map(Self::by_name).collect()
1434 }
1435
1436 #[must_use]
1438 pub fn convert_from_ids(ids: impl IntoIterator<Item = [u8; 16]>) -> Vec<Self> {
1439 ids.into_iter().map(Self::by_id).collect()
1440 }
1441
1442 #[must_use]
1447 pub fn error_result(&self, error_code: i16) -> TopicMetadata {
1448 TopicMetadata::error(error_code, self.name.as_deref(), self.topic_id)
1449 }
1450}
1451
1452pub struct MetadataRequest;
1454
1455impl MetadataRequest {
1456 #[must_use]
1462 pub const fn is_all_topics(version: i16, topics: Option<&[MetadataRequestTopic]>) -> bool {
1463 match topics {
1464 None => true,
1465 Some(topics) => topics.is_empty() && version == 0,
1466 }
1467 }
1468
1469 #[must_use]
1481 pub const fn all_topics() -> (Option<&'static [MetadataRequestTopic]>, bool) {
1482 (None, true)
1483 }
1484
1485 #[must_use]
1502 pub fn for_topic_ids(ids: Option<&[[u8; 16]]>) -> (Option<Vec<MetadataRequestTopic>>, bool) {
1503 (
1504 ids.map(|ids| MetadataRequestTopic::convert_from_ids(ids.iter().copied())),
1505 false,
1506 )
1507 }
1508
1509 #[must_use]
1523 pub fn for_topic_names(
1524 names: Option<&[String]>,
1525 allow_auto: bool,
1526 ) -> (Option<Vec<MetadataRequestTopic>>, bool) {
1527 (
1528 names.map(|names| MetadataRequestTopic::convert_from_names(names.iter().cloned())),
1529 allow_auto,
1530 )
1531 }
1532
1533 #[must_use]
1549 pub fn for_topic_names_version(
1550 names: Option<&[String]>,
1551 allow_auto: bool,
1552 allowed_version: i16,
1553 ) -> (Option<Vec<MetadataRequestTopic>>, bool, i16, i16) {
1554 Self::for_topic_names_range(names, allow_auto, allowed_version, allowed_version)
1555 }
1556
1557 #[must_use]
1573 pub fn for_topic_names_range(
1574 names: Option<&[String]>,
1575 allow_auto: bool,
1576 min_version: i16,
1577 max_version: i16,
1578 ) -> (Option<Vec<MetadataRequestTopic>>, bool, i16, i16) {
1579 let (topics, allow_auto) = Self::for_topic_names(names, allow_auto);
1580 (topics, allow_auto, min_version, max_version)
1581 }
1582
1583 #[must_use]
1588 pub fn topic_ids(version: i16, topics: Option<&[MetadataRequestTopic]>) -> Vec<[u8; 16]> {
1589 if Self::is_all_topics(version, topics) || version < 10 {
1590 Vec::new()
1591 } else {
1592 topics
1593 .unwrap_or(&[])
1594 .iter()
1595 .map(|topic| topic.topic_id)
1596 .collect()
1597 }
1598 }
1599
1600 #[must_use]
1605 pub fn topics(
1606 version: i16,
1607 topics: Option<&[MetadataRequestTopic]>,
1608 ) -> Option<Vec<Option<&str>>> {
1609 if Self::is_all_topics(version, topics) {
1610 None
1611 } else {
1612 Some(
1613 topics
1614 .unwrap_or(&[])
1615 .iter()
1616 .map(|topic| topic.name.as_deref())
1617 .collect(),
1618 )
1619 }
1620 }
1621
1622 #[must_use]
1638 pub fn error_response(
1639 topics: Option<&[MetadataRequestTopic]>,
1640 error_code: i16,
1641 ) -> MetadataResponse {
1642 MetadataResponse {
1643 throttle_time_ms: 0,
1644 brokers: Vec::new(),
1645 cluster_id: None,
1646 controller_id: MetadataResponse::NO_CONTROLLER_ID,
1647 topics: topics
1648 .unwrap_or(&[])
1649 .iter()
1650 .map(|topic| topic.error_result(error_code))
1651 .collect(),
1652 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
1653 error_code,
1654 }
1655 }
1656
1657 pub fn build(
1669 version: i16,
1670 topics: Option<&[MetadataRequestTopic]>,
1671 allow_auto: bool,
1672 ) -> Result<()> {
1673 if version < 1 {
1674 return Err(Error::Unsupported(
1675 "MetadataRequest versions older than 1 are not supported.".into(),
1676 ));
1677 }
1678 if !allow_auto && version < 4 {
1679 return Err(Error::Unsupported(
1680 "MetadataRequest versions older than 4 don't support the allowAutoTopicCreation field"
1681 .into(),
1682 ));
1683 }
1684 if let Some(topics) = topics {
1685 for t in topics {
1686 if t.name.is_none() && version < 12 {
1687 return Err(Error::Unsupported(format!(
1688 "MetadataRequest version {version} does not support null topic names."
1689 )));
1690 }
1691 if t.topic_id != [0; 16] && version < 12 {
1692 return Err(Error::Unsupported(format!(
1693 "MetadataRequest version {version} does not support non-zero topic IDs."
1694 )));
1695 }
1696 }
1697 }
1698 Ok(())
1699 }
1700}
1701
1702pub fn encode_metadata_request(
1707 buf: &mut BytesMut,
1708 version: i16,
1709 topics: Option<&[String]>,
1710 allow_auto: bool,
1711) -> crate::error::Result<()> {
1712 encode_metadata_request_with(buf, version, topics, allow_auto, false)
1713}
1714
1715pub fn encode_metadata_request_with(
1724 buf: &mut BytesMut,
1725 version: i16,
1726 topics: Option<&[String]>,
1727 allow_auto: bool,
1728 include_topic_authorized_operations: bool,
1729) -> crate::error::Result<()> {
1730 let owned = topics.map(|names| MetadataRequestTopic::convert_from_names(names.iter().cloned()));
1731 encode_metadata_request_topics(
1732 buf,
1733 version,
1734 owned.as_deref(),
1735 allow_auto,
1736 include_topic_authorized_operations,
1737 )
1738}
1739
1740pub fn encode_metadata_request_topics(
1753 buf: &mut BytesMut,
1754 version: i16,
1755 topics: Option<&[MetadataRequestTopic]>,
1756 allow_auto: bool,
1757 include_topic_authorized_operations: bool,
1758) -> crate::error::Result<()> {
1759 encode_metadata_request_topics_with_include_cluster_authorized_operations(
1760 buf,
1761 version,
1762 topics,
1763 allow_auto,
1764 include_topic_authorized_operations,
1765 false,
1766 )
1767}
1768
1769pub fn encode_metadata_request_topics_with_include_cluster_authorized_operations(
1777 buf: &mut BytesMut,
1778 version: i16,
1779 topics: Option<&[MetadataRequestTopic]>,
1780 allow_auto: bool,
1781 include_topic_authorized_operations: bool,
1782 include_cluster_authorized_operations: bool,
1783) -> crate::error::Result<()> {
1784 MetadataRequest::build(version, topics, allow_auto)?;
1785 let flexible = version >= 9;
1786 match topics {
1787 None => buf::put_array_len(buf, flexible, None)?,
1788 Some(topics) => {
1789 buf::put_array_len(buf, flexible, Some(topics.len()))?;
1790 for t in topics {
1791 if version >= 10 {
1792 buf.extend_from_slice(&t.topic_id);
1793 }
1794 buf::put_string(buf, flexible, t.name.as_deref())?;
1795 if flexible {
1796 buf::put_empty_tagged_fields(buf);
1797 }
1798 }
1799 }
1800 }
1801 if version >= 4 {
1802 buf.put_u8(u8::from(allow_auto));
1803 }
1804 if (8..=10).contains(&version) {
1805 buf.put_u8(u8::from(include_cluster_authorized_operations));
1806 }
1807 if version >= 8 {
1808 buf.put_u8(u8::from(include_topic_authorized_operations));
1809 }
1810 if flexible {
1811 buf::put_empty_tagged_fields(buf);
1812 }
1813 Ok(())
1814}
1815
1816pub fn decode_metadata_request<B: Buf>(
1823 buf: &mut B,
1824 version: i16,
1825) -> Result<(Option<Vec<String>>, bool, bool, bool)> {
1826 let (topics, allow_auto, include_topic_authorized, include_cluster_authorized) =
1827 decode_metadata_request_topics(buf, version)?;
1828 let names = topics.map(|ts| ts.into_iter().filter_map(|t| t.name).collect());
1829 Ok((
1830 names,
1831 allow_auto,
1832 include_topic_authorized,
1833 include_cluster_authorized,
1834 ))
1835}
1836
1837pub fn decode_metadata_request_topics<B: Buf>(
1845 buf: &mut B,
1846 version: i16,
1847) -> Result<(Option<Vec<MetadataRequestTopic>>, bool, bool, bool)> {
1848 let flexible = version >= 9;
1849 let topics = match buf::get_array_len(buf, flexible)? {
1850 None => None,
1851 Some(n) => {
1852 let mut topics = Vec::with_capacity(n);
1853 for _ in 0..n {
1854 let topic_id = if version >= 10 {
1855 buf::get_uuid(buf)?
1856 } else {
1857 [0; 16]
1858 };
1859 let name = buf::get_string(buf, flexible)?;
1860 if flexible {
1861 buf::skip_tagged_fields(buf)?;
1862 }
1863 topics.push(MetadataRequestTopic { name, topic_id });
1864 }
1865 Some(topics)
1866 }
1867 };
1868 let allow_auto = if version >= 4 {
1869 buf::need(buf, 1)?;
1870 buf.get_u8() != 0
1871 } else {
1872 true
1873 };
1874 let include_cluster_authorized = if (8..=10).contains(&version) {
1875 buf::need(buf, 1)?;
1876 buf.get_u8() != 0
1877 } else {
1878 false
1879 };
1880 let include_topic_authorized = if version >= 8 {
1881 buf::need(buf, 1)?;
1882 buf.get_u8() != 0
1883 } else {
1884 false
1885 };
1886 if flexible {
1887 buf::skip_tagged_fields(buf)?;
1888 }
1889 Ok((
1890 topics,
1891 allow_auto,
1892 include_topic_authorized,
1893 include_cluster_authorized,
1894 ))
1895}
1896
1897fn get_int32_array<B: Buf>(buf: &mut B, flexible: bool) -> Result<Vec<i32>> {
1898 let n = buf::get_array_len(buf, flexible)?.unwrap_or(0);
1899 let mut out = Vec::with_capacity(n);
1900 for _ in 0..n {
1901 out.push(buf::get_i32(buf)?);
1902 }
1903 Ok(out)
1904}
1905
1906fn put_int32_array(buf: &mut BytesMut, flexible: bool, items: &[i32]) -> crate::error::Result<()> {
1907 buf::put_array_len(buf, flexible, Some(items.len()))?;
1908 for v in items {
1909 buf.put_i32(*v);
1910 }
1911 Ok(())
1912}
1913
1914pub fn decode_metadata_response<B: Buf>(buf: &mut B, version: i16) -> Result<MetadataResponse> {
1916 let flexible = version >= 9;
1917 let throttle_time_ms = if version >= 3 { buf::get_i32(buf)? } else { 0 };
1918 let broker_count = buf::get_array_len(buf, flexible)?.unwrap_or(0);
1919 let mut brokers = Vec::with_capacity(broker_count);
1920 for _ in 0..broker_count {
1921 let node_id = buf::get_i32(buf)?;
1922 let host =
1923 buf::get_string(buf, flexible)?.ok_or_else(|| Error::protocol("null broker host"))?;
1924 let port = buf::get_i32(buf)?;
1925 let rack = if version >= 1 {
1926 buf::get_string(buf, flexible)?
1927 } else {
1928 None
1929 };
1930 if flexible {
1931 buf::skip_tagged_fields(buf)?;
1932 }
1933 brokers.push(Broker {
1934 node_id,
1935 host,
1936 port,
1937 rack,
1938 });
1939 }
1940 let cluster_id = if version >= 2 {
1941 buf::get_string(buf, flexible)?
1942 } else {
1943 None
1944 };
1945 let controller_id = if version >= 1 {
1946 buf::get_i32(buf)?
1947 } else {
1948 MetadataResponse::NO_CONTROLLER_ID
1949 };
1950 let topic_count = buf::get_array_len(buf, flexible)?.unwrap_or(0);
1951 let mut topics = Vec::with_capacity(topic_count);
1952 for _ in 0..topic_count {
1953 let error_code = buf::get_i16(buf)?;
1954 let name = buf::get_string(buf, flexible)?;
1955 let topic_id = if version >= 10 {
1956 buf::get_uuid(buf)?
1957 } else {
1958 [0u8; 16]
1959 };
1960 let is_internal = if version >= 1 {
1961 buf::get_bool(buf)?
1962 } else {
1963 false
1964 };
1965 let part_count = buf::get_array_len(buf, flexible)?.unwrap_or(0);
1966 let mut partitions = Vec::with_capacity(part_count);
1967 for _ in 0..part_count {
1968 let error_code = buf::get_i16(buf)?;
1969 let partition_index = buf::get_i32(buf)?;
1970 let leader_id = buf::get_i32(buf)?;
1971 let leader_epoch = if version >= 7 {
1972 buf::get_i32(buf)?
1973 } else {
1974 RecordBatch::NO_PARTITION_LEADER_EPOCH
1975 };
1976 let replica_nodes = get_int32_array(buf, flexible)?;
1977 let isr_nodes = get_int32_array(buf, flexible)?;
1978 let offline_replicas = if version >= 5 {
1979 get_int32_array(buf, flexible)?
1980 } else {
1981 Vec::new()
1982 };
1983 if flexible {
1984 buf::skip_tagged_fields(buf)?;
1985 }
1986 partitions.push(PartitionMetadata {
1987 error_code,
1988 partition_index,
1989 leader_id,
1990 leader_epoch,
1991 replica_nodes,
1992 isr_nodes,
1993 offline_replicas,
1994 });
1995 }
1996 let topic_authorized_operations = if version >= 8 {
1997 buf::get_i32(buf)?
1998 } else {
1999 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
2000 };
2001 if flexible {
2002 buf::skip_tagged_fields(buf)?;
2003 }
2004 topics.push(TopicMetadata {
2005 error_code,
2006 name,
2007 topic_id,
2008 is_internal,
2009 partitions,
2010 topic_authorized_operations,
2011 });
2012 }
2013 let cluster_authorized_operations = if (8..=10).contains(&version) {
2014 buf::get_i32(buf)?
2015 } else {
2016 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
2017 };
2018 let error_code = if version >= 13 { buf::get_i16(buf)? } else { 0 };
2019 if flexible {
2020 buf::skip_tagged_fields(buf)?;
2021 }
2022 Ok(MetadataResponse {
2023 throttle_time_ms,
2024 brokers,
2025 cluster_id,
2026 controller_id,
2027 topics,
2028 cluster_authorized_operations,
2029 error_code,
2030 })
2031}
2032
2033pub fn encode_metadata_response(
2035 buf: &mut BytesMut,
2036 version: i16,
2037 resp: &MetadataResponse,
2038) -> crate::error::Result<()> {
2039 let flexible = version >= 9;
2040 if version >= 3 {
2041 buf.put_i32(resp.throttle_time_ms);
2042 }
2043 buf::put_array_len(buf, flexible, Some(resp.brokers.len()))?;
2044 for b in &resp.brokers {
2045 buf.put_i32(b.node_id);
2046 buf::put_string(buf, flexible, Some(&b.host))?;
2047 buf.put_i32(b.port);
2048 if version >= 1 {
2049 buf::put_string(buf, flexible, b.rack.as_deref())?;
2050 }
2051 if flexible {
2052 buf::put_empty_tagged_fields(buf);
2053 }
2054 }
2055 if version >= 2 {
2056 buf::put_string(buf, flexible, resp.cluster_id.as_deref())?;
2057 }
2058 if version >= 1 {
2059 buf.put_i32(resp.controller_id);
2060 }
2061 buf::put_array_len(buf, flexible, Some(resp.topics.len()))?;
2062 for t in &resp.topics {
2063 buf.put_i16(t.error_code);
2064 buf::put_string(buf, flexible, t.name.as_deref())?;
2065 if version >= 10 {
2066 buf.extend_from_slice(&t.topic_id);
2067 }
2068 if version >= 1 {
2069 buf.put_u8(u8::from(t.is_internal));
2070 }
2071 buf::put_array_len(buf, flexible, Some(t.partitions.len()))?;
2072 for p in &t.partitions {
2073 buf.put_i16(p.error_code);
2074 buf.put_i32(p.partition_index);
2075 buf.put_i32(p.leader_id);
2076 if version >= 7 {
2077 buf.put_i32(p.leader_epoch);
2078 }
2079 put_int32_array(buf, flexible, &p.replica_nodes)?;
2080 put_int32_array(buf, flexible, &p.isr_nodes)?;
2081 if version >= 5 {
2082 put_int32_array(buf, flexible, &p.offline_replicas)?;
2083 }
2084 if flexible {
2085 buf::put_empty_tagged_fields(buf);
2086 }
2087 }
2088 if version >= 8 {
2089 buf.put_i32(t.topic_authorized_operations);
2090 }
2091 if flexible {
2092 buf::put_empty_tagged_fields(buf);
2093 }
2094 }
2095 if (8..=10).contains(&version) {
2096 buf.put_i32(resp.cluster_authorized_operations);
2097 }
2098 if version >= 13 {
2099 buf.put_i16(resp.error_code);
2100 }
2101 if flexible {
2102 buf::put_empty_tagged_fields(buf);
2103 }
2104 Ok(())
2105}
2106
2107#[derive(Debug, Clone)]
2111pub struct ProduceTopicData {
2112 pub topic: String,
2114 pub partitions: Vec<ProducePartitionData>,
2116}
2117
2118impl ProduceTopicData {
2119 #[must_use]
2127 pub fn error_result(&self, error_code: i16) -> Vec<ProducePartitionResponse> {
2128 self.partitions
2129 .iter()
2130 .map(|p| {
2131 ProducePartitionResponse::partition_response(
2132 self.topic.clone(),
2133 p.index,
2134 error_code,
2135 )
2136 })
2137 .collect()
2138 }
2139}
2140
2141#[derive(Debug, Clone)]
2143pub struct ProducePartitionData {
2144 pub index: i32,
2146 pub records: RecordBatch,
2148}
2149
2150#[derive(Debug, Clone, PartialEq, Eq)]
2157pub struct ProduceRecordError {
2158 pub batch_index: i32,
2160 pub message: Option<String>,
2162}
2163
2164impl ProduceRecordError {
2165 #[must_use]
2170 pub fn new(batch_index: i32, message: Option<String>) -> Self {
2171 Self {
2172 batch_index,
2173 message,
2174 }
2175 }
2176
2177 #[must_use]
2179 pub fn batch_index(&self) -> i32 {
2180 self.batch_index
2181 }
2182
2183 #[must_use]
2185 pub fn message(&self) -> Option<&str> {
2186 self.message.as_deref()
2187 }
2188}
2189
2190impl fmt::Display for ProduceRecordError {
2191 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2192 f.write_str("RecordError(batchIndex=")?;
2193 write!(f, "{}", self.batch_index)?;
2194 f.write_str(", message=")?;
2195 match self.message.as_deref() {
2196 None => f.write_str("null")?,
2197 Some(message) => {
2198 f.write_str("'")?;
2199 f.write_str(message)?;
2200 f.write_str("'")?;
2201 }
2202 }
2203 f.write_str(")")
2204 }
2205}
2206
2207#[derive(Debug, Clone, PartialEq, Eq)]
2239pub struct ProducePartitionResponse {
2240 pub topic: String,
2242 pub partition: i32,
2244 pub error_code: i16,
2246 pub base_offset: i64,
2248 pub log_append_time_ms: i64,
2250 pub log_start_offset: i64,
2252 pub current_leader_id: i32,
2255 pub current_leader_epoch: i32,
2258 pub record_errors: Vec<ProduceRecordError>,
2261 pub error_message: Option<String>,
2264}
2265
2266impl ProducePartitionResponse {
2267 pub const INVALID_OFFSET: i64 = -1;
2269
2270 #[must_use]
2280 pub fn partition_response(topic: impl Into<String>, partition: i32, error_code: i16) -> Self {
2281 Self::partition_response_with_offsets(
2282 topic,
2283 partition,
2284 error_code,
2285 Self::INVALID_OFFSET,
2286 RecordBatch::NO_TIMESTAMP,
2287 Self::INVALID_OFFSET,
2288 )
2289 }
2290
2291 #[must_use]
2304 pub fn partition_response_with_offsets(
2305 topic: impl Into<String>,
2306 partition: i32,
2307 error_code: i16,
2308 base_offset: i64,
2309 log_append_time_ms: i64,
2310 log_start_offset: i64,
2311 ) -> Self {
2312 Self {
2313 topic: topic.into(),
2314 partition,
2315 error_code,
2316 base_offset,
2317 log_append_time_ms,
2318 log_start_offset,
2319 current_leader_id: MetadataResponse::NO_LEADER_ID,
2320 current_leader_epoch: RecordBatch::NO_PARTITION_LEADER_EPOCH,
2321 record_errors: Vec::new(),
2322 error_message: None,
2323 }
2324 }
2325
2326 #[must_use]
2341 pub fn partition_response_with_message(
2342 topic: impl Into<String>,
2343 partition: i32,
2344 error_code: i16,
2345 error_message: Option<String>,
2346 ) -> Self {
2347 let mut part = Self::partition_response(topic, partition, error_code);
2348 part.error_message = error_message;
2349 part
2350 }
2351
2352 #[must_use]
2368 pub fn partition_response_with_record_errors(
2369 topic: impl Into<String>,
2370 partition: i32,
2371 error_code: i16,
2372 base_offset: i64,
2373 log_append_time_ms: i64,
2374 log_start_offset: i64,
2375 record_errors: Vec<ProduceRecordError>,
2376 ) -> Self {
2377 let mut part = Self::partition_response_with_offsets(
2378 topic,
2379 partition,
2380 error_code,
2381 base_offset,
2382 log_append_time_ms,
2383 log_start_offset,
2384 );
2385 part.record_errors = record_errors;
2386 part
2387 }
2388
2389 #[must_use]
2406 #[expect(
2407 clippy::too_many_arguments,
2408 reason = "Java PartitionResponse six-arg plus topic and partition stored on this type"
2409 )]
2410 pub fn partition_response_with_record_errors_and_message(
2411 topic: impl Into<String>,
2412 partition: i32,
2413 error_code: i16,
2414 base_offset: i64,
2415 log_append_time_ms: i64,
2416 log_start_offset: i64,
2417 record_errors: Vec<ProduceRecordError>,
2418 error_message: Option<String>,
2419 ) -> Self {
2420 let mut part = Self::partition_response_with_record_errors(
2421 topic,
2422 partition,
2423 error_code,
2424 base_offset,
2425 log_append_time_ms,
2426 log_start_offset,
2427 record_errors,
2428 );
2429 part.error_message = error_message;
2430 part
2431 }
2432
2433 #[must_use]
2453 #[expect(
2454 clippy::too_many_arguments,
2455 reason = "Java PartitionResponse CurrentLeader constructor plus topic and partition stored on this type"
2456 )]
2457 pub fn partition_response_with_current_leader(
2458 topic: impl Into<String>,
2459 partition: i32,
2460 error_code: i16,
2461 base_offset: i64,
2462 log_append_time_ms: i64,
2463 log_start_offset: i64,
2464 record_errors: Vec<ProduceRecordError>,
2465 error_message: Option<String>,
2466 current_leader_id: i32,
2467 current_leader_epoch: i32,
2468 ) -> Self {
2469 let mut part = Self::partition_response_with_record_errors_and_message(
2470 topic,
2471 partition,
2472 error_code,
2473 base_offset,
2474 log_append_time_ms,
2475 log_start_offset,
2476 record_errors,
2477 error_message,
2478 );
2479 part.current_leader_id = current_leader_id;
2480 part.current_leader_epoch = current_leader_epoch;
2481 part
2482 }
2483}
2484
2485pub struct ProduceRequest;
2487
2488impl ProduceRequest {
2489 pub const LAST_STABLE_VERSION_BEFORE_TRANSACTION_V2: i16 = 11;
2491
2492 #[must_use]
2494 pub const fn is_transaction_v2_requested(version: i16) -> bool {
2495 version > Self::LAST_STABLE_VERSION_BEFORE_TRANSACTION_V2
2496 }
2497
2498 #[must_use]
2511 pub const fn builder(use_transaction_v1_version: bool) -> (i16, i16) {
2512 if use_transaction_v1_version {
2513 (3, Self::LAST_STABLE_VERSION_BEFORE_TRANSACTION_V2)
2514 } else {
2515 (3, 12)
2516 }
2517 }
2518
2519 #[must_use]
2524 pub fn has_transactional_records<I, B>(partitions: I) -> bool
2525 where
2526 I: IntoIterator<Item = B>,
2527 B: AsRef<[RecordBatch]>,
2528 {
2529 partitions.into_iter().any(|part| {
2530 part.as_ref()
2531 .first()
2532 .is_some_and(RecordBatch::is_transactional)
2533 })
2534 }
2535
2536 pub fn validate_records(version: i16, batches: &[RecordBatch]) -> Result<()> {
2552 let mut iter = batches.iter();
2553 let Some(entry) = iter.next() else {
2554 return Err(Error::protocol(format!(
2555 "Produce requests with version {version} must have at least one record batch per partition"
2556 )));
2557 };
2558 if entry.magic() != RecordBatch::MAGIC_VALUE_V2 {
2559 return Err(Error::protocol(format!(
2560 "Produce requests with version {version} are only allowed to contain record batches with magic version 2"
2561 )));
2562 }
2563 if version < 7 && entry.attributes & 0x07 == 4 {
2565 return Err(Error::Unsupported(format!(
2566 "Produce requests with version {version} are not allowed to use ZStandard compression"
2567 )));
2568 }
2569 if iter.next().is_some() {
2570 return Err(Error::protocol(format!(
2571 "Produce requests with version {version} are only allowed to contain exactly one record batch per partition"
2572 )));
2573 }
2574 Ok(())
2575 }
2576
2577 pub fn build(version: i16, topics: &[ProduceTopicData]) -> Result<()> {
2587 for topic in topics {
2588 for partition in &topic.partitions {
2589 Self::validate_records(version, std::slice::from_ref(&partition.records))?;
2590 }
2591 }
2592 Ok(())
2593 }
2594
2595 pub fn partition_sizes(topics: &[ProduceTopicData]) -> Result<HashMap<(String, i32), i32>> {
2602 let mut sizes: HashMap<(String, i32), i32> = HashMap::new();
2603 for topic in topics {
2604 for partition in &topic.partitions {
2605 let size = partition.records.size_in_bytes()?;
2606 let _size = sizes
2607 .entry((topic.topic.clone(), partition.index))
2608 .and_modify(|prev| *prev = prev.wrapping_add(size))
2609 .or_insert(size);
2610 }
2611 }
2612 Ok(sizes)
2613 }
2614
2615 #[must_use]
2629 pub fn error_response(
2630 acks: i16,
2631 topics: &[ProduceTopicData],
2632 error_code: i16,
2633 ) -> Option<Vec<ProducePartitionResponse>> {
2634 if acks == 0 {
2635 return None;
2636 }
2637 let mut order: Vec<(String, i32)> = Vec::new();
2638 let mut seen: HashSet<(String, i32)> = HashSet::new();
2639 for topic in topics {
2640 for partition in &topic.partitions {
2641 let key = (topic.topic.clone(), partition.index);
2642 if !seen.insert(key.clone()) {
2643 continue;
2644 }
2645 order.push(key);
2646 }
2647 }
2648 Some(
2649 order
2650 .into_iter()
2651 .map(|(topic, partition)| {
2652 ProducePartitionResponse::partition_response(topic, partition, error_code)
2653 })
2654 .collect(),
2655 )
2656 }
2657
2658 #[must_use]
2669 pub fn error_counts(topics: &[ProduceTopicData], error_code: i16) -> HashMap<i16, i32> {
2670 let mut seen: HashSet<(String, i32)> = HashSet::new();
2671 let mut n = 0i32;
2672 for topic in topics {
2673 for partition in &topic.partitions {
2674 if seen.insert((topic.topic.clone(), partition.index)) {
2675 n = n.wrapping_add(1);
2676 }
2677 }
2678 }
2679 HashMap::from([(error_code, n)])
2680 }
2681}
2682
2683pub struct ProduceResponse;
2685
2686impl ProduceResponse {
2687 #[must_use]
2689 pub const fn should_client_throttle(version: i16) -> bool {
2690 version >= 6
2691 }
2692
2693 #[must_use]
2697 pub fn error_counts(partitions: &[ProducePartitionResponse]) -> HashMap<i16, i32> {
2698 let mut counts = HashMap::new();
2699 for partition in partitions {
2700 let count = counts.entry(partition.error_code).or_insert(0);
2701 *count += 1;
2702 }
2703 counts
2704 }
2705
2706 #[must_use]
2719 pub fn to_data(
2720 partitions: &[ProducePartitionResponse],
2721 ) -> Vec<(String, Vec<ProducePartitionResponse>)> {
2722 let mut topics = Vec::<(String, Vec<ProducePartitionResponse>)>::new();
2723 for partition in partitions {
2724 if let Some((_, parts)) = topics.iter_mut().find(|(name, _)| name == &partition.topic) {
2725 parts.push(partition.clone());
2726 } else {
2727 topics.push((partition.topic.clone(), vec![partition.clone()]));
2728 }
2729 }
2730 topics
2731 }
2732}
2733
2734fn produce_flexible(version: i16) -> Result<bool> {
2744 match version {
2745 3..=8 => Ok(false),
2746 9..=12 => Ok(true),
2747 other => Err(Error::protocol(format!(
2748 "Produce version {other} is not implemented"
2749 ))),
2750 }
2751}
2752
2753fn encode_current_leader(leader_id: i32, leader_epoch: i32) -> Bytes {
2754 let mut inner = BytesMut::new();
2755 inner.put_i32(leader_id);
2756 inner.put_i32(leader_epoch);
2757 buf::put_empty_tagged_fields(&mut inner);
2758 inner.freeze()
2759}
2760
2761fn decode_current_leader(value: &Bytes) -> Result<(i32, i32)> {
2762 let mut cur = value.as_ref();
2763 let leader_id = buf::get_i32(&mut cur)?;
2764 let leader_epoch = buf::get_i32(&mut cur)?;
2765 buf::skip_tagged_fields(&mut cur)?;
2766 leftover_empty(&cur, "CurrentLeader")?;
2767 Ok((leader_id, leader_epoch))
2768}
2769
2770fn encode_produce_partition_tags(
2771 buf: &mut BytesMut,
2772 version: i16,
2773 current_leader_id: i32,
2774 current_leader_epoch: i32,
2775) -> Result<()> {
2776 if version >= 10 && current_leader_id >= 0 {
2777 buf::put_tagged_fields(
2778 buf,
2779 &[(
2780 0,
2781 encode_current_leader(current_leader_id, current_leader_epoch),
2782 )],
2783 )
2784 } else {
2785 buf::put_empty_tagged_fields(buf);
2786 Ok(())
2787 }
2788}
2789
2790fn decode_produce_partition_tags<B: Buf>(buf: &mut B, version: i16) -> Result<(i32, i32)> {
2791 let tags = buf::get_tagged_fields(buf)?;
2792 let mut current_leader_id = MetadataResponse::NO_LEADER_ID;
2793 let mut current_leader_epoch = RecordBatch::NO_PARTITION_LEADER_EPOCH;
2794 if version >= 10 {
2795 for (tag, value) in tags {
2796 if tag == 0 {
2797 (current_leader_id, current_leader_epoch) = decode_current_leader(&value)?;
2798 }
2799 }
2800 }
2801 Ok((current_leader_id, current_leader_epoch))
2802}
2803
2804pub fn encode_produce_request(
2806 buf: &mut BytesMut,
2807 version: i16,
2808 transactional_id: Option<&str>,
2809 acks: i16,
2810 timeout_ms: i32,
2811 topics: &[ProduceTopicData],
2812) -> Result<()> {
2813 let flexible = produce_flexible(version)?;
2814 if version >= 3 {
2815 buf::put_string(buf, flexible, transactional_id)?;
2816 }
2817 buf.put_i16(acks);
2818 buf.put_i32(timeout_ms);
2819 buf::put_array_len(buf, flexible, Some(topics.len()))?;
2820 for topic in topics {
2821 if version <= 12 {
2822 buf::put_string(buf, flexible, Some(&topic.topic))?;
2823 } else {
2824 buf.extend_from_slice(&[0u8; 16]);
2825 }
2826 buf::put_array_len(buf, flexible, Some(topic.partitions.len()))?;
2827 for part in &topic.partitions {
2828 buf.put_i32(part.index);
2829 if flexible {
2830 let mut recs = BytesMut::new();
2831 records::encode_record_batch(&mut recs, &part.records)?;
2832 buf::put_bytes(buf, flexible, Some(&recs))?;
2833 buf::put_empty_tagged_fields(buf);
2834 } else {
2835 let len_pos = buf.len();
2836 buf.put_i32(0);
2837 records::encode_record_batch(buf, &part.records)?;
2838 let rec_len = buf::i32_from_usize(buf.len().saturating_sub(len_pos + 4))?;
2839 buf::patch_i32(buf, len_pos, rec_len)?;
2840 }
2841 }
2842 if flexible {
2843 buf::put_empty_tagged_fields(buf);
2844 }
2845 }
2846 if flexible {
2847 buf::put_empty_tagged_fields(buf);
2848 }
2849 Ok(())
2850}
2851
2852pub fn decode_produce_request<B: Buf>(
2854 buf: &mut B,
2855 version: i16,
2856) -> Result<(Option<String>, i16, i32, Vec<ProduceTopicData>)> {
2857 let flexible = produce_flexible(version)?;
2858 let transactional_id = if version >= 3 {
2859 buf::get_string(buf, flexible)?
2860 } else {
2861 None
2862 };
2863 let acks = buf::get_i16(buf)?;
2864 let timeout_ms = buf::get_i32(buf)?;
2865 let topic_count = buf::get_array_len(buf, flexible)?.unwrap_or(0);
2866 let mut topics = Vec::with_capacity(topic_count);
2867 for _ in 0..topic_count {
2868 let topic = if version <= 12 {
2869 buf::get_string(buf, flexible)?.unwrap_or_default()
2870 } else {
2871 let _id = buf::get_uuid(buf)?;
2872 String::new()
2873 };
2874 let part_count = buf::get_array_len(buf, flexible)?.unwrap_or(0);
2875 let mut partitions = Vec::with_capacity(part_count);
2876 for _ in 0..part_count {
2877 let index = buf::get_i32(buf)?;
2878 let rec_bytes = buf::get_bytes(buf, flexible)?.unwrap_or_default();
2879 let mut rec_buf = &rec_bytes[..];
2880 let records = if rec_buf.is_empty() {
2881 RecordBatch::from_records(vec![])
2882 } else {
2883 records::decode_record_batch(&mut rec_buf)?
2884 };
2885 if flexible {
2886 buf::skip_tagged_fields(buf)?;
2887 }
2888 partitions.push(ProducePartitionData { index, records });
2889 }
2890 if flexible {
2891 buf::skip_tagged_fields(buf)?;
2892 }
2893 topics.push(ProduceTopicData { topic, partitions });
2894 }
2895 if flexible {
2896 buf::skip_tagged_fields(buf)?;
2897 }
2898 Ok((transactional_id, acks, timeout_ms, topics))
2899}
2900
2901fn encode_produce_record_errors(
2904 buf: &mut BytesMut,
2905 flexible: bool,
2906 errors: &[ProduceRecordError],
2907) -> Result<()> {
2908 buf::put_array_len(buf, flexible, Some(errors.len()))?;
2909 for err in errors {
2910 buf.put_i32(err.batch_index);
2911 buf::put_string(buf, flexible, err.message.as_deref())?;
2912 if flexible {
2913 buf::put_empty_tagged_fields(buf);
2914 }
2915 }
2916 Ok(())
2917}
2918
2919fn decode_produce_record_errors<B: Buf>(
2920 buf: &mut B,
2921 flexible: bool,
2922) -> Result<Vec<ProduceRecordError>> {
2923 let n = buf::get_array_len(buf, flexible)?.unwrap_or(0);
2924 let mut errors = Vec::with_capacity(n);
2925 for _ in 0..n {
2926 let batch_index = buf::get_i32(buf)?;
2927 let message = buf::get_string(buf, flexible)?;
2928 if flexible {
2929 buf::skip_tagged_fields(buf)?;
2930 }
2931 errors.push(ProduceRecordError {
2932 batch_index,
2933 message,
2934 });
2935 }
2936 Ok(errors)
2937}
2938
2939pub fn encode_produce_response(
2943 buf: &mut BytesMut,
2944 version: i16,
2945 parts: &[ProducePartitionResponse],
2946) -> crate::error::Result<()> {
2947 encode_produce_response_with_endpoints(buf, version, parts, &[])
2948}
2949
2950pub fn encode_produce_response_with_throttle(
2960 buf: &mut BytesMut,
2961 version: i16,
2962 parts: &[ProducePartitionResponse],
2963 throttle_time_ms: i32,
2964) -> crate::error::Result<()> {
2965 encode_produce_response_fields(buf, version, parts, &[], throttle_time_ms)
2966}
2967
2968pub fn encode_produce_response_with_endpoints(
2972 buf: &mut BytesMut,
2973 version: i16,
2974 parts: &[ProducePartitionResponse],
2975 endpoints: &[NodeEndpoint],
2976) -> crate::error::Result<()> {
2977 encode_produce_response_fields(buf, version, parts, endpoints, 0)
2978}
2979
2980fn encode_produce_response_fields(
2981 buf: &mut BytesMut,
2982 version: i16,
2983 parts: &[ProducePartitionResponse],
2984 endpoints: &[NodeEndpoint],
2985 throttle_time_ms: i32,
2986) -> crate::error::Result<()> {
2987 let flexible = produce_flexible(version)?;
2988 let mut order: Vec<String> = Vec::new();
2990 for p in parts {
2991 if !order.iter().any(|t| t == &p.topic) {
2992 order.push(p.topic.clone());
2993 }
2994 }
2995 buf::put_array_len(buf, flexible, Some(order.len()))?;
2996 for topic in &order {
2997 if version <= 12 {
2998 buf::put_string(buf, flexible, Some(topic))?;
2999 } else {
3000 buf.extend_from_slice(&[0u8; 16]);
3001 }
3002 let grouped: Vec<&ProducePartitionResponse> =
3003 parts.iter().filter(|p| &p.topic == topic).collect();
3004 buf::put_array_len(buf, flexible, Some(grouped.len()))?;
3005 for p in grouped {
3006 buf.put_i32(p.partition);
3007 buf.put_i16(p.error_code);
3008 buf.put_i64(p.base_offset);
3009 if version >= 2 {
3010 buf.put_i64(p.log_append_time_ms);
3011 }
3012 if version >= 5 {
3013 buf.put_i64(p.log_start_offset);
3014 }
3015 if version >= 8 {
3016 encode_produce_record_errors(buf, flexible, &p.record_errors)?;
3017 buf::put_string(buf, flexible, p.error_message.as_deref())?;
3018 }
3019 if flexible {
3020 encode_produce_partition_tags(
3021 buf,
3022 version,
3023 p.current_leader_id,
3024 p.current_leader_epoch,
3025 )?;
3026 }
3027 }
3028 if flexible {
3029 buf::put_empty_tagged_fields(buf);
3030 }
3031 }
3032 if version >= 1 {
3033 buf.put_i32(throttle_time_ms);
3034 }
3035 if flexible {
3036 encode_top_level_node_endpoints(buf, version >= 10, endpoints)?;
3037 }
3038 Ok(())
3039}
3040
3041pub fn decode_produce_response<B: Buf>(
3050 buf: &mut B,
3051 version: i16,
3052) -> Result<(Vec<ProducePartitionResponse>, Vec<NodeEndpoint>, i32)> {
3053 let flexible = produce_flexible(version)?;
3054 let topic_count = buf::get_array_len(buf, flexible)?.unwrap_or(0);
3055 let mut out = Vec::new();
3056 for _ in 0..topic_count {
3057 let topic = if version <= 12 {
3058 buf::get_string(buf, flexible)?.unwrap_or_default()
3059 } else {
3060 let _id = buf::get_uuid(buf)?;
3061 String::new()
3062 };
3063 let part_count = buf::get_array_len(buf, flexible)?.unwrap_or(0);
3064 for _ in 0..part_count {
3065 let partition = buf::get_i32(buf)?;
3066 let error_code = buf::get_i16(buf)?;
3067 let base_offset = buf::get_i64(buf)?;
3068 let log_append_time_ms = if version >= 2 {
3069 buf::get_i64(buf)?
3070 } else {
3071 RecordBatch::NO_TIMESTAMP
3072 };
3073 let log_start_offset = if version >= 5 {
3074 buf::get_i64(buf)?
3075 } else {
3076 ProducePartitionResponse::INVALID_OFFSET
3077 };
3078 let (record_errors, error_message) = if version >= 8 {
3079 (
3080 decode_produce_record_errors(buf, flexible)?,
3081 buf::get_string(buf, flexible)?,
3082 )
3083 } else {
3084 (Vec::new(), None)
3085 };
3086 let (current_leader_id, current_leader_epoch) = if flexible {
3087 decode_produce_partition_tags(buf, version)?
3088 } else {
3089 (
3090 MetadataResponse::NO_LEADER_ID,
3091 RecordBatch::NO_PARTITION_LEADER_EPOCH,
3092 )
3093 };
3094 out.push(ProducePartitionResponse {
3095 topic: topic.clone(),
3096 partition,
3097 error_code,
3098 base_offset,
3099 log_append_time_ms,
3100 log_start_offset,
3101 current_leader_id,
3102 current_leader_epoch,
3103 record_errors,
3104 error_message,
3105 });
3106 }
3107 if flexible {
3108 buf::skip_tagged_fields(buf)?;
3109 }
3110 }
3111 let throttle_time_ms = if version >= 1 { buf::get_i32(buf)? } else { 0 };
3112 let endpoints = if flexible {
3113 decode_top_level_node_endpoints(buf, version >= 10)?
3114 } else {
3115 Vec::new()
3116 };
3117 Ok((out, endpoints, throttle_time_ms))
3118}
3119
3120#[cfg(test)]
3121mod tests {
3122 use super::*;
3123 use crate::protocol::records::{Compression, Record};
3124 use bytes::Bytes;
3125 use std::collections::{HashMap, HashSet};
3126
3127 #[test]
3128 fn produce_transaction_v2_version_cap_matches_java() {
3129 assert_eq!(
3130 ProduceRequest::LAST_STABLE_VERSION_BEFORE_TRANSACTION_V2,
3131 11
3132 );
3133 assert!(!ProduceRequest::is_transaction_v2_requested(11));
3134 assert!(ProduceRequest::is_transaction_v2_requested(12));
3135 assert!(!ProduceRequest::is_transaction_v2_requested(10));
3136 assert!(!ProduceResponse::should_client_throttle(5));
3137 assert!(ProduceResponse::should_client_throttle(6));
3138 assert!(ProduceResponse::error_counts(&[]).is_empty());
3139 let counts = ProduceResponse::error_counts(&[
3140 ProducePartitionResponse::partition_response("t", 0, 0),
3141 ProducePartitionResponse::partition_response(
3142 "t",
3143 1,
3144 crate::error::NOT_LEADER_OR_FOLLOWER,
3145 ),
3146 ProducePartitionResponse::partition_response("t", 2, 0),
3147 ProducePartitionResponse::partition_response(
3148 "u",
3149 0,
3150 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
3151 ),
3152 ]);
3153 assert_eq!(
3154 counts,
3155 HashMap::from([
3156 (0, 2),
3157 (crate::error::NOT_LEADER_OR_FOLLOWER, 1),
3158 (crate::error::UNKNOWN_TOPIC_OR_PARTITION, 1),
3159 ])
3160 );
3161 }
3162
3163 #[test]
3164 fn has_transactional_records_matches_java_request_utils() {
3165 let rec = Record {
3166 offset: 0,
3167 timestamp: 1,
3168 key: None,
3169 value: Some(Bytes::from_static(b"x")),
3170 headers: vec![],
3171 };
3172 let plain = RecordBatch::from_records(vec![rec.clone()]);
3173 let tx = RecordBatch::from_records(vec![rec]).with_transactional(true);
3174 assert!(!ProduceRequest::has_transactional_records(
3175 std::iter::empty::<&[RecordBatch]>()
3176 ));
3177 assert!(!ProduceRequest::has_transactional_records([&[][..]]));
3178 assert!(!ProduceRequest::has_transactional_records([
3179 std::slice::from_ref(&plain)
3180 ]));
3181 assert!(ProduceRequest::has_transactional_records([
3182 std::slice::from_ref(&tx)
3183 ]));
3184 let later_tx = [plain.clone(), tx.clone()];
3185 assert!(
3186 !ProduceRequest::has_transactional_records([&later_tx[..]]),
3187 "Java inspects only the first batch of each partition"
3188 );
3189 let first_tx = [tx.clone(), plain.clone()];
3190 assert!(ProduceRequest::has_transactional_records([&first_tx[..]]));
3191 assert!(ProduceRequest::has_transactional_records([
3192 &[][..],
3193 std::slice::from_ref(&tx)
3194 ]));
3195 }
3196
3197 #[test]
3198 fn produce_request_validate_records_matches_java() {
3199 let rec = Record {
3208 offset: 0,
3209 timestamp: 1,
3210 key: None,
3211 value: Some(Bytes::from_static(b"x")),
3212 headers: vec![],
3213 };
3214 let none = RecordBatch::from_records(vec![rec.clone()]);
3215 let gzip = none.clone().with_compression(Compression::Gzip);
3216 let snappy = none.clone().with_compression(Compression::Snappy);
3217 let lz4 = none.clone().with_compression(Compression::Lz4);
3218 let mut zstd = none.clone();
3219 zstd.attributes = (zstd.attributes & !0x07) | 4;
3220 let two = [none.clone(), gzip.clone()];
3221 let two_zstd = [zstd.clone(), zstd.clone()];
3222 for version in 3..=12_i16 {
3223 let empty = ProduceRequest::validate_records(version, &[]).unwrap_err();
3224 assert!(
3225 matches!(empty, Error::Protocol(_)),
3226 "empty batches is Java InvalidRecordException, got {empty}"
3227 );
3228 assert!(
3229 empty
3230 .to_string()
3231 .contains("must have at least one record batch per partition"),
3232 "got {empty}"
3233 );
3234 let many = ProduceRequest::validate_records(version, &two).unwrap_err();
3235 assert!(
3236 matches!(many, Error::Protocol(_)),
3237 "two batches is Java InvalidRecordException, got {many}"
3238 );
3239 assert!(
3240 many.to_string()
3241 .contains("exactly one record batch per partition"),
3242 "got {many}"
3243 );
3244 ProduceRequest::validate_records(version, std::slice::from_ref(&none)).unwrap();
3245 ProduceRequest::validate_records(version, std::slice::from_ref(&gzip)).unwrap();
3246 ProduceRequest::validate_records(version, std::slice::from_ref(&snappy)).unwrap();
3247 ProduceRequest::validate_records(version, std::slice::from_ref(&lz4)).unwrap();
3248 assert_eq!(none.magic(), RecordBatch::MAGIC_VALUE_V2);
3249 if version < 7 {
3250 let zstd_err =
3251 ProduceRequest::validate_records(version, std::slice::from_ref(&zstd))
3252 .unwrap_err();
3253 assert!(
3254 matches!(zstd_err, Error::Unsupported(_)),
3255 "zstd below v7 is Java UnsupportedCompressionTypeException, got {zstd_err}"
3256 );
3257 assert!(
3258 zstd_err
3259 .to_string()
3260 .contains("are not allowed to use ZStandard compression"),
3261 "got {zstd_err}"
3262 );
3263 let zstd_first = ProduceRequest::validate_records(version, &two_zstd).unwrap_err();
3264 assert!(
3265 matches!(zstd_first, Error::Unsupported(_)),
3266 "zstd is checked before the exactly-one-batch rule, got {zstd_first}"
3267 );
3268 } else {
3269 ProduceRequest::validate_records(version, std::slice::from_ref(&zstd)).unwrap();
3270 }
3271 }
3272 leftover_empty_produce_validate_records(3, &none);
3273 leftover_empty_produce_validate_records(9, &none);
3274 leftover_empty_produce_validate_records(3, &gzip);
3275 leftover_empty_produce_validate_records(9, &gzip);
3276 }
3277
3278 fn leftover_empty_produce_validate_records(version: i16, batch: &RecordBatch) {
3279 ProduceRequest::validate_records(version, std::slice::from_ref(batch)).unwrap();
3280 let topics = [ProduceTopicData {
3281 topic: "t".into(),
3282 partitions: vec![ProducePartitionData {
3283 index: 0,
3284 records: batch.clone(),
3285 }],
3286 }];
3287 let mut buf = BytesMut::new();
3288 encode_produce_request(&mut buf, version, None, 1, 0, &topics).unwrap();
3289 let mut cur = buf.as_ref();
3290 let (.., decoded_topics) = decode_produce_request(&mut cur, version).unwrap();
3291 leftover_empty(
3292 &cur,
3293 match version {
3294 3 => "Produce v3 validateRecords leftover-empty",
3295 9 => "Produce v9 validateRecords leftover-empty",
3296 _ => "Produce validateRecords leftover-empty",
3297 },
3298 )
3299 .unwrap();
3300 assert_eq!(
3301 decoded_topics.len(),
3302 1,
3303 "v{version} validateRecords encode kept one topic"
3304 );
3305 buf.clear();
3306 encode_produce_request(&mut buf, version, None, 1, 0, &[]).unwrap();
3307 let mut cur = buf.as_ref();
3308 let (.., decoded_empty) = decode_produce_request(&mut cur, version).unwrap();
3309 leftover_empty(
3310 &cur,
3311 match version {
3312 3 => "Produce v3 validateRecords empty leftover-empty",
3313 9 => "Produce v9 validateRecords empty leftover-empty",
3314 _ => "Produce validateRecords empty leftover-empty",
3315 },
3316 )
3317 .unwrap();
3318 assert!(decoded_empty.is_empty());
3319 }
3320
3321 #[test]
3322 fn produce_request_build_matches_java() {
3323 let rec = Record {
3330 offset: 0,
3331 timestamp: 1,
3332 key: None,
3333 value: Some(Bytes::from_static(b"x")),
3334 headers: vec![],
3335 };
3336 let none = RecordBatch::from_records(vec![rec]);
3337 let mut zstd = none.clone();
3338 zstd.attributes = (zstd.attributes & !0x07) | 4;
3339 let ok = [ProduceTopicData {
3340 topic: "t".into(),
3341 partitions: vec![
3342 ProducePartitionData {
3343 index: 0,
3344 records: none.clone(),
3345 },
3346 ProducePartitionData {
3347 index: 1,
3348 records: none.clone().with_compression(Compression::Gzip),
3349 },
3350 ],
3351 }];
3352 let zstd_topics = [ProduceTopicData {
3353 topic: "t".into(),
3354 partitions: vec![ProducePartitionData {
3355 index: 0,
3356 records: zstd,
3357 }],
3358 }];
3359 for version in 3..=12_i16 {
3360 ProduceRequest::build(version, &[]).unwrap();
3361 ProduceRequest::build(version, &ok).unwrap();
3362 if version < 7 {
3363 let zstd_err = ProduceRequest::build(version, &zstd_topics).unwrap_err();
3364 assert!(
3365 matches!(zstd_err, Error::Unsupported(_)),
3366 "zstd below v7 is Java UnsupportedCompressionTypeException, got {zstd_err}"
3367 );
3368 assert!(
3369 zstd_err
3370 .to_string()
3371 .contains("are not allowed to use ZStandard compression"),
3372 "got {zstd_err}"
3373 );
3374 } else {
3375 ProduceRequest::build(version, &zstd_topics).unwrap();
3376 }
3377 }
3378 leftover_empty_produce_build(3, &ok);
3379 leftover_empty_produce_build(9, &ok);
3380 leftover_empty_produce_build(3, &[]);
3381 leftover_empty_produce_build(9, &[]);
3382 }
3383
3384 fn leftover_empty_produce_build(version: i16, topics: &[ProduceTopicData]) {
3385 ProduceRequest::build(version, topics).unwrap();
3386 let mut buf = BytesMut::new();
3387 encode_produce_request(&mut buf, version, None, 1, 0, topics).unwrap();
3388 let mut cur = buf.as_ref();
3389 let (.., decoded) = decode_produce_request(&mut cur, version).unwrap();
3390 leftover_empty(
3391 &cur,
3392 match (version, topics.is_empty()) {
3393 (3, false) => "Produce v3 Builder.build leftover-empty",
3394 (9, false) => "Produce v9 Builder.build leftover-empty",
3395 (3, true) => "Produce v3 Builder.build empty leftover-empty",
3396 (9, true) => "Produce v9 Builder.build empty leftover-empty",
3397 (_, false) => "Produce Builder.build leftover-empty",
3398 (_, true) => "Produce Builder.build empty leftover-empty",
3399 },
3400 )
3401 .unwrap();
3402 assert_eq!(decoded.len(), topics.len());
3403 }
3404
3405 #[test]
3406 fn produce_request_builder_matches_java() {
3407 let (oldest, latest) = ProduceRequest::builder(false);
3416 assert_eq!(oldest, 3);
3417 assert_eq!(latest, 12);
3418 assert_eq!(
3419 ProduceRequest::builder(true),
3420 (3, ProduceRequest::LAST_STABLE_VERSION_BEFORE_TRANSACTION_V2)
3421 );
3422 assert_eq!(ProduceRequest::builder(true), (3, 11));
3423 assert!(!ProduceRequest::is_transaction_v2_requested(
3424 ProduceRequest::builder(true).1
3425 ));
3426 assert!(ProduceRequest::is_transaction_v2_requested(latest));
3427 let rec = Record {
3428 offset: 0,
3429 timestamp: 1,
3430 key: None,
3431 value: Some(Bytes::from_static(b"x")),
3432 headers: vec![],
3433 };
3434 let topics = [ProduceTopicData {
3435 topic: "t".into(),
3436 partitions: vec![ProducePartitionData {
3437 index: 0,
3438 records: RecordBatch::from_records(vec![rec]),
3439 }],
3440 }];
3441 leftover_empty_produce_builder(oldest, &topics);
3442 leftover_empty_produce_builder(oldest, &[]);
3443 leftover_empty_produce_builder(11, &topics);
3444 leftover_empty_produce_builder(11, &[]);
3445 leftover_empty_produce_builder(latest, &topics);
3446 leftover_empty_produce_builder(latest, &[]);
3447 }
3448
3449 fn leftover_empty_produce_builder(version: i16, topics: &[ProduceTopicData]) {
3450 let mut buf = BytesMut::new();
3451 encode_produce_request(&mut buf, version, None, 1, 0, topics).unwrap();
3452 let mut cur = buf.as_ref();
3453 let (.., decoded) = decode_produce_request(&mut cur, version).unwrap();
3454 leftover_empty(
3455 &cur,
3456 match (version, topics.is_empty()) {
3457 (3, false) => "Produce v3 builder leftover-empty",
3458 (11, false) => "Produce v11 builder leftover-empty",
3459 (12, false) => "Produce v12 builder leftover-empty",
3460 (3, true) => "Produce v3 builder empty leftover-empty",
3461 (11, true) => "Produce v11 builder empty leftover-empty",
3462 (12, true) => "Produce v12 builder empty leftover-empty",
3463 (_, false) => "Produce builder leftover-empty",
3464 (_, true) => "Produce builder empty leftover-empty",
3465 },
3466 )
3467 .unwrap();
3468 assert_eq!(decoded.len(), topics.len());
3469 }
3470
3471 #[test]
3472 fn produce_partition_sizes_matches_java() {
3473 assert!(ProduceRequest::partition_sizes(&[]).unwrap().is_empty());
3477 let rec = Record {
3478 offset: 0,
3479 timestamp: 1,
3480 key: None,
3481 value: Some(Bytes::from_static(b"x")),
3482 headers: vec![],
3483 };
3484 let batch = RecordBatch::from_records(vec![rec]);
3485 let one_size = batch.size_in_bytes().unwrap();
3486 let one = [ProduceTopicData {
3487 topic: "t".into(),
3488 partitions: vec![
3489 ProducePartitionData {
3490 index: 0,
3491 records: batch.clone(),
3492 },
3493 ProducePartitionData {
3494 index: 3,
3495 records: batch.clone(),
3496 },
3497 ],
3498 }];
3499 assert_eq!(
3500 ProduceRequest::partition_sizes(&one).unwrap(),
3501 HashMap::from([(("t".into(), 0), one_size), (("t".into(), 3), one_size)])
3502 );
3503 let dup = [
3504 ProduceTopicData {
3505 topic: "a".into(),
3506 partitions: vec![ProducePartitionData {
3507 index: 0,
3508 records: batch.clone(),
3509 }],
3510 },
3511 ProduceTopicData {
3512 topic: "a".into(),
3513 partitions: vec![
3514 ProducePartitionData {
3515 index: 0,
3516 records: batch.clone(),
3517 },
3518 ProducePartitionData {
3519 index: 1,
3520 records: batch.clone(),
3521 },
3522 ],
3523 },
3524 ];
3525 assert_eq!(
3526 ProduceRequest::partition_sizes(&dup).unwrap(),
3527 HashMap::from([
3528 (("a".into(), 0), one_size.wrapping_add(one_size)),
3529 (("a".into(), 1), one_size),
3530 ])
3531 );
3532 let mut buf = BytesMut::new();
3533 encode_produce_request(&mut buf, 3, None, 1, 1000, &dup).unwrap();
3534 let mut cur = buf.as_ref();
3535 let decoded = decode_produce_request(&mut cur, 3).unwrap().3;
3536 leftover_empty(&cur, "Produce v3 partitionSizes").unwrap();
3537 assert_eq!(
3538 ProduceRequest::partition_sizes(&decoded).unwrap(),
3539 ProduceRequest::partition_sizes(&dup).unwrap()
3540 );
3541 buf.clear();
3542 encode_produce_request(&mut buf, 9, None, 1, 1000, &dup).unwrap();
3543 let mut cur = buf.as_ref();
3544 let decoded = decode_produce_request(&mut cur, 9).unwrap().3;
3545 leftover_empty(&cur, "Produce v9 partitionSizes").unwrap();
3546 assert_eq!(
3547 ProduceRequest::partition_sizes(&decoded).unwrap(),
3548 ProduceRequest::partition_sizes(&dup).unwrap()
3549 );
3550 }
3551
3552 #[test]
3553 fn produce_response_to_data_matches_java() {
3554 assert!(ProduceResponse::to_data(&[]).is_empty());
3558 let a0 = ProducePartitionResponse::partition_response("a", 0, 0);
3559 let b0 =
3560 ProducePartitionResponse::partition_response("b", 0, crate::error::CORRUPT_MESSAGE);
3561 let a1 = ProducePartitionResponse::partition_response(
3562 "a",
3563 1,
3564 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
3565 );
3566 let a0_dup = ProducePartitionResponse::partition_response("a", 0, 0);
3567 let grouped =
3568 ProduceResponse::to_data(&[a0.clone(), b0.clone(), a1.clone(), a0_dup.clone()]);
3569 assert_eq!(grouped.len(), 2);
3570 assert_eq!(grouped[0].0, "a");
3571 assert_eq!(grouped[0].1, vec![a0.clone(), a1.clone(), a0_dup.clone()]);
3572 assert_eq!(grouped[1].0, "b");
3573 assert_eq!(grouped[1].1, vec![b0.clone()]);
3574
3575 let parts = [a0.clone(), b0.clone(), a1.clone(), a0_dup.clone()];
3576 for version in [3_i16, 9] {
3577 let grouped = ProduceResponse::to_data(&parts);
3578 assert_eq!(grouped.len(), 2);
3579 let mut buf = BytesMut::new();
3580 encode_produce_response(&mut buf, version, &parts).unwrap();
3581 let mut cur = buf.as_ref();
3582 let (decoded, endpoints, ..) = decode_produce_response(&mut cur, version).unwrap();
3583 assert!(endpoints.is_empty());
3584 assert_eq!(ProduceResponse::to_data(&decoded).len(), 2);
3585 leftover_empty(
3586 &cur,
3587 match version {
3588 3 => "Produce v3 toData leftover-empty",
3589 9 => "Produce v9 toData leftover-empty",
3590 _ => "Produce toData leftover-empty",
3591 },
3592 )
3593 .unwrap();
3594 }
3595 for version in [3_i16, 9] {
3596 assert!(ProduceResponse::to_data(&[]).is_empty());
3597 let mut buf = BytesMut::new();
3598 encode_produce_response(&mut buf, version, &[]).unwrap();
3599 let mut cur = buf.as_ref();
3600 let (decoded, endpoints, ..) = decode_produce_response(&mut cur, version).unwrap();
3601 assert!(decoded.is_empty());
3602 assert!(endpoints.is_empty());
3603 leftover_empty(
3604 &cur,
3605 match version {
3606 3 => "Produce v3 toData empty leftover-empty",
3607 9 => "Produce v9 toData empty leftover-empty",
3608 _ => "Produce toData empty leftover-empty",
3609 },
3610 )
3611 .unwrap();
3612 }
3613 }
3614
3615 #[test]
3616 fn produce_response_throttle_time_ms_matches_java() {
3617 let parts: Vec<ProducePartitionResponse> = vec![];
3629 for version in 3..=12 {
3630 let mut buf = BytesMut::new();
3631 encode_produce_response_with_throttle(&mut buf, version, &parts, 3_600_000).unwrap();
3632 let mut cur = buf.as_ref();
3633 let (decoded, endpoints, throttle) =
3634 decode_produce_response(&mut cur, version).unwrap();
3635 assert_eq!(decoded, parts);
3636 assert!(endpoints.is_empty());
3637 assert_eq!(throttle, 3_600_000);
3638 assert!(
3639 cur.is_empty(),
3640 "Produce v{version} ThrottleTimeMs leftover-empty"
3641 );
3642 }
3643
3644 let mut with = BytesMut::new();
3645 encode_produce_response_with_throttle(&mut with, 3, &parts, 3_600_000).unwrap();
3646 let mut zero = BytesMut::new();
3647 encode_produce_response_with_throttle(&mut zero, 3, &parts, 0).unwrap();
3648 assert_ne!(
3649 &with[..],
3650 &zero[..],
3651 "v3 ThrottleTimeMs is not always the JSON default 0"
3652 );
3653 let mut conv = BytesMut::new();
3654 encode_produce_response(&mut conv, 3, &parts).unwrap();
3655 assert_eq!(
3656 &conv[..],
3657 &zero[..],
3658 "encode_produce_response still writes ThrottleTimeMs 0"
3659 );
3660 assert_eq!(
3661 &with[..4],
3662 &[0, 0, 0, 0],
3663 "empty-Responses ThrottleTimeMs comes after the topic array"
3664 );
3665
3666 let mut v8_with = BytesMut::new();
3667 encode_produce_response_with_throttle(&mut v8_with, 8, &parts, 3_600_000).unwrap();
3668 for version in 4..=8 {
3669 let mut buf = BytesMut::new();
3670 encode_produce_response_with_throttle(&mut buf, version, &parts, 3_600_000).unwrap();
3671 assert_eq!(
3672 &with[..],
3673 &buf[..],
3674 "empty-Responses ThrottleTimeMs bodies: v3 == v{version}"
3675 );
3676 }
3677 let mut v9_with = BytesMut::new();
3678 encode_produce_response_with_throttle(&mut v9_with, 9, &parts, 3_600_000).unwrap();
3679 assert_ne!(
3680 &v8_with[..],
3681 &v9_with[..],
3682 "v9 adds compact tagged fields after ThrottleTimeMs"
3683 );
3684 for version in 10..=12 {
3685 let mut buf = BytesMut::new();
3686 encode_produce_response_with_throttle(&mut buf, version, &parts, 3_600_000).unwrap();
3687 assert_eq!(
3688 &v9_with[..],
3689 &buf[..],
3690 "empty-Responses ThrottleTimeMs bodies: v9 == v{version}"
3691 );
3692 }
3693 }
3694
3695 #[test]
3696 fn produce_error_response_matches_java() {
3697 let rec = Record {
3701 offset: 0,
3702 timestamp: 1,
3703 key: None,
3704 value: Some(Bytes::from_static(b"x")),
3705 headers: vec![],
3706 };
3707 let batch = RecordBatch::from_records(vec![rec]);
3708 let dup = [
3709 ProduceTopicData {
3710 topic: "a".into(),
3711 partitions: vec![ProducePartitionData {
3712 index: 0,
3713 records: batch.clone(),
3714 }],
3715 },
3716 ProduceTopicData {
3717 topic: "a".into(),
3718 partitions: vec![
3719 ProducePartitionData {
3720 index: 0,
3721 records: batch.clone(),
3722 },
3723 ProducePartitionData {
3724 index: 1,
3725 records: batch,
3726 },
3727 ],
3728 },
3729 ];
3730 assert!(
3731 ProduceRequest::error_response(0, &dup, crate::error::CORRUPT_MESSAGE).is_none(),
3732 "acks 0 is Java null"
3733 );
3734 assert!(ProduceRequest::error_response(0, &[], 0).is_none());
3735 let empty = ProduceRequest::error_response(1, &[], crate::error::CORRUPT_MESSAGE)
3736 .expect("acks 1 empty is empty Topics");
3737 assert!(empty.is_empty());
3738 let parts =
3739 ProduceRequest::error_response(1, &dup, crate::error::UNKNOWN_TOPIC_OR_PARTITION)
3740 .expect("acks 1");
3741 assert_eq!(
3742 parts,
3743 vec![
3744 ProducePartitionResponse::partition_response(
3745 "a",
3746 0,
3747 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
3748 ),
3749 ProducePartitionResponse::partition_response(
3750 "a",
3751 1,
3752 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
3753 ),
3754 ]
3755 );
3756 assert_eq!(
3757 parts.len(),
3758 2,
3759 "duplicate (topic, partition) is one partitionSizes key"
3760 );
3761 let per_request: Vec<_> = dup.iter().flat_map(|t| t.error_result(0)).collect();
3762 assert_eq!(per_request.len(), 3, "error_result keeps every partition");
3763 leftover_empty_produce_error_response(3, &parts, &empty);
3764 leftover_empty_produce_error_response(9, &parts, &empty);
3765 }
3766
3767 #[test]
3768 fn produce_request_error_counts_matches_java() {
3769 let rec = Record {
3774 offset: 0,
3775 timestamp: 1,
3776 key: None,
3777 value: Some(Bytes::from_static(b"x")),
3778 headers: vec![],
3779 };
3780 let batch = RecordBatch::from_records(vec![rec]);
3781 let dup = [
3782 ProduceTopicData {
3783 topic: "a".into(),
3784 partitions: vec![ProducePartitionData {
3785 index: 0,
3786 records: batch.clone(),
3787 }],
3788 },
3789 ProduceTopicData {
3790 topic: "a".into(),
3791 partitions: vec![
3792 ProducePartitionData {
3793 index: 0,
3794 records: batch.clone(),
3795 },
3796 ProducePartitionData {
3797 index: 1,
3798 records: batch,
3799 },
3800 ],
3801 },
3802 ];
3803 assert_eq!(
3804 ProduceRequest::error_counts(&[], crate::error::CORRUPT_MESSAGE),
3805 HashMap::from([(crate::error::CORRUPT_MESSAGE, 0)]),
3806 "empty Topics is a singleton 0, not an empty map"
3807 );
3808 assert_eq!(
3809 ProduceRequest::error_counts(&dup, crate::error::UNKNOWN_TOPIC_OR_PARTITION),
3810 HashMap::from([(crate::error::UNKNOWN_TOPIC_OR_PARTITION, 2)]),
3811 "duplicate (topic, partition) is one partitionSizes key"
3812 );
3813 assert_eq!(
3814 ProduceRequest::error_counts(&dup, crate::error::CORRUPT_MESSAGE)
3815 .get(&crate::error::CORRUPT_MESSAGE)
3816 .copied(),
3817 Some(2),
3818 "acks is not consulted; getErrorResponse would be None for acks 0"
3819 );
3820 assert_eq!(
3821 ProduceRequest::error_counts(&dup, 0),
3822 HashMap::from([(0, 2)]),
3823 "NONE is still a map key"
3824 );
3825 assert!(
3826 ProduceResponse::error_counts(&[]).is_empty(),
3827 "response errorCounts on empty is empty, unlike request errorCounts"
3828 );
3829
3830 let mut buf = BytesMut::new();
3831 encode_produce_request(&mut buf, 3, None, 0, 1000, &dup).unwrap();
3832 let mut cur = buf.as_ref();
3833 let decoded = decode_produce_request(&mut cur, 3).unwrap().3;
3834 leftover_empty(&cur, "Produce v3 Request.errorCounts").unwrap();
3835 assert_eq!(
3836 ProduceRequest::error_counts(&decoded, crate::error::CORRUPT_MESSAGE),
3837 ProduceRequest::error_counts(&dup, crate::error::CORRUPT_MESSAGE)
3838 );
3839 buf.clear();
3840 encode_produce_request(&mut buf, 9, None, 0, 1000, &dup).unwrap();
3841 let mut cur = buf.as_ref();
3842 let decoded = decode_produce_request(&mut cur, 9).unwrap().3;
3843 leftover_empty(&cur, "Produce v9 Request.errorCounts").unwrap();
3844 assert_eq!(
3845 ProduceRequest::error_counts(&decoded, crate::error::CORRUPT_MESSAGE),
3846 ProduceRequest::error_counts(&dup, crate::error::CORRUPT_MESSAGE)
3847 );
3848 buf.clear();
3849 encode_produce_request(&mut buf, 3, None, 1, 1000, &[]).unwrap();
3850 let mut cur = buf.as_ref();
3851 let decoded = decode_produce_request(&mut cur, 3).unwrap().3;
3852 leftover_empty(&cur, "Produce v3 Request.errorCounts empty").unwrap();
3853 assert_eq!(
3854 ProduceRequest::error_counts(&decoded, crate::error::CORRUPT_MESSAGE),
3855 HashMap::from([(crate::error::CORRUPT_MESSAGE, 0)])
3856 );
3857 }
3858
3859 fn leftover_empty_produce_error_response(
3860 version: i16,
3861 parts: &[ProducePartitionResponse],
3862 empty: &[ProducePartitionResponse],
3863 ) {
3864 let mut buf = BytesMut::new();
3865 encode_produce_response(&mut buf, version, parts).unwrap();
3866 let mut cur = buf.as_ref();
3867 let (decoded, endpoints, ..) = decode_produce_response(&mut cur, version).unwrap();
3868 assert!(endpoints.is_empty());
3869 assert_eq!(decoded, parts);
3870 leftover_empty(
3871 &cur,
3872 match version {
3873 3 => "Produce v3 getErrorResponse from partitionSizes",
3874 9 => "Produce v9 getErrorResponse from partitionSizes",
3875 _ => "Produce getErrorResponse from partitionSizes",
3876 },
3877 )
3878 .unwrap();
3879 buf.clear();
3880 encode_produce_response(&mut buf, version, empty).unwrap();
3881 let mut cur = buf.as_ref();
3882 let (decoded, endpoints, ..) = decode_produce_response(&mut cur, version).unwrap();
3883 assert!(endpoints.is_empty());
3884 assert_eq!(decoded, empty);
3885 leftover_empty(
3886 &cur,
3887 match version {
3888 3 => "Produce v3 getErrorResponse from partitionSizes empty",
3889 9 => "Produce v9 getErrorResponse from partitionSizes empty",
3890 _ => "Produce getErrorResponse from partitionSizes empty",
3891 },
3892 )
3893 .unwrap();
3894 }
3895
3896 #[test]
3897 fn produce_partition_response_matches_java() {
3898 assert_eq!(ProducePartitionResponse::INVALID_OFFSET, -1);
3899 assert_eq!(RecordBatch::NO_TIMESTAMP, -1);
3900 assert_eq!(MetadataResponse::NO_LEADER_ID, -1);
3901 assert_eq!(RecordBatch::NO_PARTITION_LEADER_EPOCH, -1);
3902 let none = ProducePartitionResponse::partition_response("t", 0, 0);
3903 assert_eq!(none.topic, "t");
3904 assert_eq!(none.partition, 0);
3905 assert_eq!(none.error_code, 0);
3906 assert_eq!(none.base_offset, ProducePartitionResponse::INVALID_OFFSET);
3907 assert_eq!(none.log_append_time_ms, RecordBatch::NO_TIMESTAMP);
3908 assert_eq!(
3909 none.log_start_offset,
3910 ProducePartitionResponse::INVALID_OFFSET
3911 );
3912 assert_eq!(none.current_leader_id, MetadataResponse::NO_LEADER_ID);
3913 assert_eq!(
3914 none.current_leader_epoch,
3915 RecordBatch::NO_PARTITION_LEADER_EPOCH
3916 );
3917 assert!(none.record_errors.is_empty());
3918 assert!(none.error_message.is_none());
3919 let unknown = ProducePartitionResponse::partition_response(
3920 "missing",
3921 3,
3922 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
3923 );
3924 assert_eq!(unknown.topic, "missing");
3925 assert_eq!(unknown.partition, 3);
3926 assert_eq!(unknown.error_code, crate::error::UNKNOWN_TOPIC_OR_PARTITION);
3927 assert_eq!(
3928 unknown.base_offset,
3929 ProducePartitionResponse::INVALID_OFFSET
3930 );
3931 assert_eq!(unknown.log_append_time_ms, RecordBatch::NO_TIMESTAMP);
3932 assert_eq!(
3933 unknown.log_start_offset,
3934 ProducePartitionResponse::INVALID_OFFSET
3935 );
3936 let topic = ProduceTopicData {
3937 topic: "t".into(),
3938 partitions: vec![
3939 ProducePartitionData {
3940 index: 0,
3941 records: RecordBatch::from_records(vec![]),
3942 },
3943 ProducePartitionData {
3944 index: 3,
3945 records: RecordBatch::from_records(vec![]),
3946 },
3947 ],
3948 };
3949 let result = topic.error_result(crate::error::UNKNOWN_TOPIC_OR_PARTITION);
3950 assert_eq!(
3951 result,
3952 vec![
3953 ProducePartitionResponse::partition_response(
3954 "t",
3955 0,
3956 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
3957 ),
3958 ProducePartitionResponse::partition_response(
3959 "t",
3960 3,
3961 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
3962 ),
3963 ]
3964 );
3965 let mut buf = BytesMut::new();
3966 encode_produce_response(&mut buf, 8, &result).unwrap();
3967 let mut cur = buf.as_ref();
3968 let (decoded, endpoints, ..) = decode_produce_response(&mut cur, 8).unwrap();
3969 assert!(endpoints.is_empty());
3970 assert_eq!(decoded, result);
3971 leftover_empty(&cur, "Produce getErrorResponse v8").unwrap();
3972 }
3973
3974 #[test]
3975 fn produce_partition_response_with_offsets_matches_java() {
3976 assert_eq!(
3984 ProducePartitionResponse::partition_response("t", 0, 0),
3985 ProducePartitionResponse::partition_response_with_offsets(
3986 "t",
3987 0,
3988 0,
3989 ProducePartitionResponse::INVALID_OFFSET,
3990 RecordBatch::NO_TIMESTAMP,
3991 ProducePartitionResponse::INVALID_OFFSET,
3992 )
3993 );
3994 let with = ProducePartitionResponse::partition_response_with_offsets(
3995 "t",
3996 1,
3997 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
3998 7,
3999 99,
4000 3,
4001 );
4002 assert_eq!(with.topic, "t");
4003 assert_eq!(with.partition, 1);
4004 assert_eq!(with.error_code, crate::error::UNKNOWN_TOPIC_OR_PARTITION);
4005 assert_eq!(with.base_offset, 7);
4006 assert_eq!(with.log_append_time_ms, 99);
4007 assert_eq!(with.log_start_offset, 3);
4008 assert_eq!(with.current_leader_id, MetadataResponse::NO_LEADER_ID);
4009 assert_eq!(
4010 with.current_leader_epoch,
4011 RecordBatch::NO_PARTITION_LEADER_EPOCH
4012 );
4013 assert!(with.record_errors.is_empty());
4014 assert!(with.error_message.is_none());
4015 leftover_produce_partition_response_with_offsets(3, std::slice::from_ref(&with));
4016 leftover_produce_partition_response_with_offsets(3, &[]);
4017 leftover_produce_partition_response_with_offsets(5, std::slice::from_ref(&with));
4018 leftover_produce_partition_response_with_offsets(5, &[]);
4019 leftover_produce_partition_response_with_offsets(8, std::slice::from_ref(&with));
4020 leftover_produce_partition_response_with_offsets(8, &[]);
4021 leftover_produce_partition_response_with_offsets(12, std::slice::from_ref(&with));
4022 leftover_produce_partition_response_with_offsets(12, &[]);
4023 }
4024
4025 fn leftover_produce_partition_response_with_offsets(
4026 version: i16,
4027 parts: &[ProducePartitionResponse],
4028 ) {
4029 let mut buf = BytesMut::new();
4030 encode_produce_response(&mut buf, version, parts).unwrap();
4031 let mut cur = buf.as_ref();
4032 let (decoded, endpoints, throttle) = decode_produce_response(&mut cur, version).unwrap();
4033 assert!(endpoints.is_empty());
4034 assert_eq!(throttle, 0);
4035 if version < 5 {
4036 assert_eq!(decoded.len(), parts.len());
4037 for (got, want) in decoded.iter().zip(parts) {
4038 assert_eq!(got.topic, want.topic);
4039 assert_eq!(got.partition, want.partition);
4040 assert_eq!(got.error_code, want.error_code);
4041 assert_eq!(got.base_offset, want.base_offset);
4042 assert_eq!(got.log_append_time_ms, want.log_append_time_ms);
4043 assert_eq!(
4044 got.log_start_offset,
4045 ProducePartitionResponse::INVALID_OFFSET,
4046 "Produce v{version} omits LogStartOffset; decode fills INVALID_OFFSET"
4047 );
4048 assert_eq!(got.current_leader_id, MetadataResponse::NO_LEADER_ID);
4049 assert_eq!(
4050 got.current_leader_epoch,
4051 RecordBatch::NO_PARTITION_LEADER_EPOCH
4052 );
4053 assert!(got.record_errors.is_empty());
4054 assert!(got.error_message.is_none());
4055 }
4056 } else {
4057 assert_eq!(decoded, parts);
4058 }
4059 leftover_empty(
4060 &cur,
4061 match (version, parts.is_empty()) {
4062 (3, false) => "Produce v3 PartitionResponse.offsets leftover-empty",
4063 (3, true) => "Produce v3 PartitionResponse.offsets empty leftover-empty",
4064 (5, false) => "Produce v5 PartitionResponse.offsets leftover-empty",
4065 (5, true) => "Produce v5 PartitionResponse.offsets empty leftover-empty",
4066 (8, false) => "Produce v8 PartitionResponse.offsets leftover-empty",
4067 (8, true) => "Produce v8 PartitionResponse.offsets empty leftover-empty",
4068 (12, false) => "Produce v12 PartitionResponse.offsets leftover-empty",
4069 (12, true) => "Produce v12 PartitionResponse.offsets empty leftover-empty",
4070 _ => "Produce PartitionResponse.offsets leftover-empty",
4071 },
4072 )
4073 .unwrap();
4074 }
4075
4076 #[test]
4077 fn produce_partition_response_with_message_matches_java() {
4078 assert_eq!(
4087 ProducePartitionResponse::partition_response("t", 0, 0),
4088 ProducePartitionResponse::partition_response_with_message("t", 0, 0, None)
4089 );
4090 let empty_msg = ProducePartitionResponse::partition_response_with_message(
4091 "t",
4092 0,
4093 0,
4094 Some(String::new()),
4095 );
4096 assert_eq!(empty_msg.error_message.as_deref(), Some(""));
4097 assert_ne!(
4098 empty_msg,
4099 ProducePartitionResponse::partition_response("t", 0, 0)
4100 );
4101 let with = ProducePartitionResponse::partition_response_with_message(
4102 "t",
4103 1,
4104 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
4105 Some("unknown topic".into()),
4106 );
4107 assert_eq!(with.topic, "t");
4108 assert_eq!(with.partition, 1);
4109 assert_eq!(with.error_code, crate::error::UNKNOWN_TOPIC_OR_PARTITION);
4110 assert_eq!(with.base_offset, ProducePartitionResponse::INVALID_OFFSET);
4111 assert_eq!(with.log_append_time_ms, RecordBatch::NO_TIMESTAMP);
4112 assert_eq!(
4113 with.log_start_offset,
4114 ProducePartitionResponse::INVALID_OFFSET
4115 );
4116 assert_eq!(with.current_leader_id, MetadataResponse::NO_LEADER_ID);
4117 assert_eq!(
4118 with.current_leader_epoch,
4119 RecordBatch::NO_PARTITION_LEADER_EPOCH
4120 );
4121 assert!(with.record_errors.is_empty());
4122 assert_eq!(with.error_message.as_deref(), Some("unknown topic"));
4123 leftover_produce_partition_response_with_message(3, std::slice::from_ref(&with));
4124 leftover_produce_partition_response_with_message(3, &[]);
4125 leftover_produce_partition_response_with_message(5, std::slice::from_ref(&with));
4126 leftover_produce_partition_response_with_message(5, &[]);
4127 leftover_produce_partition_response_with_message(8, std::slice::from_ref(&with));
4128 leftover_produce_partition_response_with_message(8, &[]);
4129 leftover_produce_partition_response_with_message(12, std::slice::from_ref(&with));
4130 leftover_produce_partition_response_with_message(12, &[]);
4131 }
4132
4133 fn leftover_produce_partition_response_with_message(
4134 version: i16,
4135 parts: &[ProducePartitionResponse],
4136 ) {
4137 let mut buf = BytesMut::new();
4138 encode_produce_response(&mut buf, version, parts).unwrap();
4139 let mut cur = buf.as_ref();
4140 let (decoded, endpoints, throttle) = decode_produce_response(&mut cur, version).unwrap();
4141 assert!(endpoints.is_empty());
4142 assert_eq!(throttle, 0);
4143 if version >= 8 {
4144 assert_eq!(decoded, parts);
4145 } else {
4146 assert_eq!(decoded.len(), parts.len());
4147 for (got, want) in decoded.iter().zip(parts) {
4148 assert_eq!(got.topic, want.topic);
4149 assert_eq!(got.partition, want.partition);
4150 assert_eq!(got.error_code, want.error_code);
4151 assert_eq!(got.base_offset, want.base_offset);
4152 assert_eq!(got.log_append_time_ms, want.log_append_time_ms);
4153 if version < 5 {
4154 assert_eq!(
4155 got.log_start_offset,
4156 ProducePartitionResponse::INVALID_OFFSET,
4157 "Produce v{version} omits LogStartOffset; decode fills INVALID_OFFSET"
4158 );
4159 } else {
4160 assert_eq!(got.log_start_offset, want.log_start_offset);
4161 }
4162 assert_eq!(got.current_leader_id, MetadataResponse::NO_LEADER_ID);
4163 assert_eq!(
4164 got.current_leader_epoch,
4165 RecordBatch::NO_PARTITION_LEADER_EPOCH
4166 );
4167 assert!(got.record_errors.is_empty());
4168 assert!(
4169 got.error_message.is_none(),
4170 "Produce v{version} omits ErrorMessage; decode fills null"
4171 );
4172 }
4173 }
4174 leftover_empty(
4175 &cur,
4176 match (version, parts.is_empty()) {
4177 (3, false) => "Produce v3 PartitionResponse.errorMessage leftover-empty",
4178 (3, true) => "Produce v3 PartitionResponse.errorMessage empty leftover-empty",
4179 (5, false) => "Produce v5 PartitionResponse.errorMessage leftover-empty",
4180 (5, true) => "Produce v5 PartitionResponse.errorMessage empty leftover-empty",
4181 (8, false) => "Produce v8 PartitionResponse.errorMessage leftover-empty",
4182 (8, true) => "Produce v8 PartitionResponse.errorMessage empty leftover-empty",
4183 (12, false) => "Produce v12 PartitionResponse.errorMessage leftover-empty",
4184 (12, true) => "Produce v12 PartitionResponse.errorMessage empty leftover-empty",
4185 _ => "Produce PartitionResponse.errorMessage leftover-empty",
4186 },
4187 )
4188 .unwrap();
4189 }
4190
4191 #[test]
4192 fn produce_partition_response_with_record_errors_matches_java() {
4193 assert_eq!(
4203 ProducePartitionResponse::partition_response_with_offsets(
4204 "t",
4205 0,
4206 0,
4207 ProducePartitionResponse::INVALID_OFFSET,
4208 RecordBatch::NO_TIMESTAMP,
4209 ProducePartitionResponse::INVALID_OFFSET,
4210 ),
4211 ProducePartitionResponse::partition_response_with_record_errors(
4212 "t",
4213 0,
4214 0,
4215 ProducePartitionResponse::INVALID_OFFSET,
4216 RecordBatch::NO_TIMESTAMP,
4217 ProducePartitionResponse::INVALID_OFFSET,
4218 Vec::new(),
4219 )
4220 );
4221 let with = ProducePartitionResponse::partition_response_with_record_errors(
4222 "t",
4223 1,
4224 crate::error::INVALID_RECORD,
4225 7,
4226 99,
4227 3,
4228 vec![
4229 ProduceRecordError::new(0, None),
4230 ProduceRecordError::new(0, Some("dup".into())),
4231 ProduceRecordError::new(2, Some("bad".into())),
4232 ],
4233 );
4234 assert_eq!(with.topic, "t");
4235 assert_eq!(with.partition, 1);
4236 assert_eq!(with.error_code, crate::error::INVALID_RECORD);
4237 assert_eq!(with.base_offset, 7);
4238 assert_eq!(with.log_append_time_ms, 99);
4239 assert_eq!(with.log_start_offset, 3);
4240 assert_eq!(with.current_leader_id, MetadataResponse::NO_LEADER_ID);
4241 assert_eq!(
4242 with.current_leader_epoch,
4243 RecordBatch::NO_PARTITION_LEADER_EPOCH
4244 );
4245 assert_eq!(
4246 with.record_errors,
4247 vec![
4248 ProduceRecordError::new(0, None),
4249 ProduceRecordError::new(0, Some("dup".into())),
4250 ProduceRecordError::new(2, Some("bad".into())),
4251 ]
4252 );
4253 assert!(with.error_message.is_none());
4254 leftover_produce_partition_response_with_record_errors(3, std::slice::from_ref(&with));
4255 leftover_produce_partition_response_with_record_errors(3, &[]);
4256 leftover_produce_partition_response_with_record_errors(5, std::slice::from_ref(&with));
4257 leftover_produce_partition_response_with_record_errors(5, &[]);
4258 leftover_produce_partition_response_with_record_errors(8, std::slice::from_ref(&with));
4259 leftover_produce_partition_response_with_record_errors(8, &[]);
4260 leftover_produce_partition_response_with_record_errors(12, std::slice::from_ref(&with));
4261 leftover_produce_partition_response_with_record_errors(12, &[]);
4262 }
4263
4264 fn leftover_produce_partition_response_with_record_errors(
4265 version: i16,
4266 parts: &[ProducePartitionResponse],
4267 ) {
4268 let mut buf = BytesMut::new();
4269 encode_produce_response(&mut buf, version, parts).unwrap();
4270 let mut cur = buf.as_ref();
4271 let (decoded, endpoints, throttle) = decode_produce_response(&mut cur, version).unwrap();
4272 assert!(endpoints.is_empty());
4273 assert_eq!(throttle, 0);
4274 if version >= 8 {
4275 assert_eq!(decoded, parts);
4276 } else {
4277 assert_eq!(decoded.len(), parts.len());
4278 for (got, want) in decoded.iter().zip(parts) {
4279 assert_eq!(got.topic, want.topic);
4280 assert_eq!(got.partition, want.partition);
4281 assert_eq!(got.error_code, want.error_code);
4282 assert_eq!(got.base_offset, want.base_offset);
4283 assert_eq!(got.log_append_time_ms, want.log_append_time_ms);
4284 if version < 5 {
4285 assert_eq!(
4286 got.log_start_offset,
4287 ProducePartitionResponse::INVALID_OFFSET,
4288 "Produce v{version} omits LogStartOffset; decode fills INVALID_OFFSET"
4289 );
4290 } else {
4291 assert_eq!(got.log_start_offset, want.log_start_offset);
4292 }
4293 assert_eq!(got.current_leader_id, MetadataResponse::NO_LEADER_ID);
4294 assert_eq!(
4295 got.current_leader_epoch,
4296 RecordBatch::NO_PARTITION_LEADER_EPOCH
4297 );
4298 assert!(
4299 got.record_errors.is_empty(),
4300 "Produce v{version} omits RecordErrors; decode fills empty"
4301 );
4302 assert!(got.error_message.is_none());
4303 }
4304 }
4305 leftover_empty(
4306 &cur,
4307 match (version, parts.is_empty()) {
4308 (3, false) => "Produce v3 PartitionResponse.recordErrors leftover-empty",
4309 (3, true) => "Produce v3 PartitionResponse.recordErrors empty leftover-empty",
4310 (5, false) => "Produce v5 PartitionResponse.recordErrors leftover-empty",
4311 (5, true) => "Produce v5 PartitionResponse.recordErrors empty leftover-empty",
4312 (8, false) => "Produce v8 PartitionResponse.recordErrors leftover-empty",
4313 (8, true) => "Produce v8 PartitionResponse.recordErrors empty leftover-empty",
4314 (12, false) => "Produce v12 PartitionResponse.recordErrors leftover-empty",
4315 (12, true) => "Produce v12 PartitionResponse.recordErrors empty leftover-empty",
4316 _ => "Produce PartitionResponse.recordErrors leftover-empty",
4317 },
4318 )
4319 .unwrap();
4320 }
4321
4322 #[test]
4323 fn produce_partition_response_with_record_errors_and_message_matches_java() {
4324 assert_eq!(
4336 ProducePartitionResponse::partition_response_with_record_errors(
4337 "t",
4338 0,
4339 0,
4340 ProducePartitionResponse::INVALID_OFFSET,
4341 RecordBatch::NO_TIMESTAMP,
4342 ProducePartitionResponse::INVALID_OFFSET,
4343 Vec::new(),
4344 ),
4345 ProducePartitionResponse::partition_response_with_record_errors_and_message(
4346 "t",
4347 0,
4348 0,
4349 ProducePartitionResponse::INVALID_OFFSET,
4350 RecordBatch::NO_TIMESTAMP,
4351 ProducePartitionResponse::INVALID_OFFSET,
4352 Vec::new(),
4353 None,
4354 )
4355 );
4356 let with = ProducePartitionResponse::partition_response_with_record_errors_and_message(
4357 "t",
4358 1,
4359 crate::error::INVALID_RECORD,
4360 7,
4361 99,
4362 3,
4363 vec![
4364 ProduceRecordError::new(0, None),
4365 ProduceRecordError::new(0, Some("dup".into())),
4366 ProduceRecordError::new(2, Some("bad".into())),
4367 ],
4368 Some("batch dropped".into()),
4369 );
4370 assert_eq!(with.topic, "t");
4371 assert_eq!(with.partition, 1);
4372 assert_eq!(with.error_code, crate::error::INVALID_RECORD);
4373 assert_eq!(with.base_offset, 7);
4374 assert_eq!(with.log_append_time_ms, 99);
4375 assert_eq!(with.log_start_offset, 3);
4376 assert_eq!(with.current_leader_id, MetadataResponse::NO_LEADER_ID);
4377 assert_eq!(
4378 with.current_leader_epoch,
4379 RecordBatch::NO_PARTITION_LEADER_EPOCH
4380 );
4381 assert_eq!(
4382 with.record_errors,
4383 vec![
4384 ProduceRecordError::new(0, None),
4385 ProduceRecordError::new(0, Some("dup".into())),
4386 ProduceRecordError::new(2, Some("bad".into())),
4387 ]
4388 );
4389 assert_eq!(with.error_message.as_deref(), Some("batch dropped"));
4390 leftover_produce_partition_response_with_record_errors_and_message(
4391 3,
4392 std::slice::from_ref(&with),
4393 );
4394 leftover_produce_partition_response_with_record_errors_and_message(3, &[]);
4395 leftover_produce_partition_response_with_record_errors_and_message(
4396 5,
4397 std::slice::from_ref(&with),
4398 );
4399 leftover_produce_partition_response_with_record_errors_and_message(5, &[]);
4400 leftover_produce_partition_response_with_record_errors_and_message(
4401 8,
4402 std::slice::from_ref(&with),
4403 );
4404 leftover_produce_partition_response_with_record_errors_and_message(8, &[]);
4405 leftover_produce_partition_response_with_record_errors_and_message(
4406 12,
4407 std::slice::from_ref(&with),
4408 );
4409 leftover_produce_partition_response_with_record_errors_and_message(12, &[]);
4410 }
4411
4412 fn leftover_produce_partition_response_with_record_errors_and_message(
4413 version: i16,
4414 parts: &[ProducePartitionResponse],
4415 ) {
4416 let mut buf = BytesMut::new();
4417 encode_produce_response(&mut buf, version, parts).unwrap();
4418 let mut cur = buf.as_ref();
4419 let (decoded, endpoints, throttle) = decode_produce_response(&mut cur, version).unwrap();
4420 assert!(endpoints.is_empty());
4421 assert_eq!(throttle, 0);
4422 if version >= 8 {
4423 assert_eq!(decoded, parts);
4424 } else {
4425 assert_eq!(decoded.len(), parts.len());
4426 for (got, want) in decoded.iter().zip(parts) {
4427 assert_eq!(got.topic, want.topic);
4428 assert_eq!(got.partition, want.partition);
4429 assert_eq!(got.error_code, want.error_code);
4430 assert_eq!(got.base_offset, want.base_offset);
4431 assert_eq!(got.log_append_time_ms, want.log_append_time_ms);
4432 if version < 5 {
4433 assert_eq!(
4434 got.log_start_offset,
4435 ProducePartitionResponse::INVALID_OFFSET,
4436 "Produce v{version} omits LogStartOffset; decode fills INVALID_OFFSET"
4437 );
4438 } else {
4439 assert_eq!(got.log_start_offset, want.log_start_offset);
4440 }
4441 assert_eq!(got.current_leader_id, MetadataResponse::NO_LEADER_ID);
4442 assert_eq!(
4443 got.current_leader_epoch,
4444 RecordBatch::NO_PARTITION_LEADER_EPOCH
4445 );
4446 assert!(
4447 got.record_errors.is_empty(),
4448 "Produce v{version} omits RecordErrors; decode fills empty"
4449 );
4450 assert!(
4451 got.error_message.is_none(),
4452 "Produce v{version} omits ErrorMessage; decode fills null"
4453 );
4454 }
4455 }
4456 leftover_empty(
4457 &cur,
4458 match (version, parts.is_empty()) {
4459 (3, false) => "Produce v3 PartitionResponse.recordErrorsAndMessage leftover-empty",
4460 (3, true) => {
4461 "Produce v3 PartitionResponse.recordErrorsAndMessage empty leftover-empty"
4462 }
4463 (5, false) => "Produce v5 PartitionResponse.recordErrorsAndMessage leftover-empty",
4464 (5, true) => {
4465 "Produce v5 PartitionResponse.recordErrorsAndMessage empty leftover-empty"
4466 }
4467 (8, false) => "Produce v8 PartitionResponse.recordErrorsAndMessage leftover-empty",
4468 (8, true) => {
4469 "Produce v8 PartitionResponse.recordErrorsAndMessage empty leftover-empty"
4470 }
4471 (12, false) => {
4472 "Produce v12 PartitionResponse.recordErrorsAndMessage leftover-empty"
4473 }
4474 (12, true) => {
4475 "Produce v12 PartitionResponse.recordErrorsAndMessage empty leftover-empty"
4476 }
4477 _ => "Produce PartitionResponse.recordErrorsAndMessage leftover-empty",
4478 },
4479 )
4480 .unwrap();
4481 }
4482
4483 #[test]
4484 fn produce_partition_response_with_current_leader_matches_java() {
4485 assert_eq!(
4498 ProducePartitionResponse::partition_response_with_record_errors_and_message(
4499 "t",
4500 0,
4501 0,
4502 ProducePartitionResponse::INVALID_OFFSET,
4503 RecordBatch::NO_TIMESTAMP,
4504 ProducePartitionResponse::INVALID_OFFSET,
4505 Vec::new(),
4506 None,
4507 ),
4508 ProducePartitionResponse::partition_response_with_current_leader(
4509 "t",
4510 0,
4511 0,
4512 ProducePartitionResponse::INVALID_OFFSET,
4513 RecordBatch::NO_TIMESTAMP,
4514 ProducePartitionResponse::INVALID_OFFSET,
4515 Vec::new(),
4516 None,
4517 MetadataResponse::NO_LEADER_ID,
4518 RecordBatch::NO_PARTITION_LEADER_EPOCH,
4519 )
4520 );
4521 let with = ProducePartitionResponse::partition_response_with_current_leader(
4522 "t",
4523 1,
4524 crate::error::NOT_LEADER_OR_FOLLOWER,
4525 7,
4526 99,
4527 3,
4528 vec![
4529 ProduceRecordError::new(0, None),
4530 ProduceRecordError::new(0, Some("dup".into())),
4531 ProduceRecordError::new(2, Some("bad".into())),
4532 ],
4533 Some("not leader".into()),
4534 2,
4535 7,
4536 );
4537 assert_eq!(with.topic, "t");
4538 assert_eq!(with.partition, 1);
4539 assert_eq!(with.error_code, crate::error::NOT_LEADER_OR_FOLLOWER);
4540 assert_eq!(with.base_offset, 7);
4541 assert_eq!(with.log_append_time_ms, 99);
4542 assert_eq!(with.log_start_offset, 3);
4543 assert_eq!(with.current_leader_id, 2);
4544 assert_eq!(with.current_leader_epoch, 7);
4545 assert_eq!(
4546 with.record_errors,
4547 vec![
4548 ProduceRecordError::new(0, None),
4549 ProduceRecordError::new(0, Some("dup".into())),
4550 ProduceRecordError::new(2, Some("bad".into())),
4551 ]
4552 );
4553 assert_eq!(with.error_message.as_deref(), Some("not leader"));
4554 leftover_produce_partition_response_with_current_leader(3, std::slice::from_ref(&with));
4555 leftover_produce_partition_response_with_current_leader(3, &[]);
4556 leftover_produce_partition_response_with_current_leader(5, std::slice::from_ref(&with));
4557 leftover_produce_partition_response_with_current_leader(5, &[]);
4558 leftover_produce_partition_response_with_current_leader(8, std::slice::from_ref(&with));
4559 leftover_produce_partition_response_with_current_leader(8, &[]);
4560 leftover_produce_partition_response_with_current_leader(12, std::slice::from_ref(&with));
4561 leftover_produce_partition_response_with_current_leader(12, &[]);
4562 }
4563
4564 fn leftover_produce_partition_response_with_current_leader(
4565 version: i16,
4566 parts: &[ProducePartitionResponse],
4567 ) {
4568 let mut buf = BytesMut::new();
4569 encode_produce_response(&mut buf, version, parts).unwrap();
4570 let mut cur = buf.as_ref();
4571 let (decoded, endpoints, throttle) = decode_produce_response(&mut cur, version).unwrap();
4572 assert!(endpoints.is_empty());
4573 assert_eq!(throttle, 0);
4574 if version >= 10 {
4575 assert_eq!(decoded, parts);
4576 } else {
4577 assert_eq!(decoded.len(), parts.len());
4578 for (got, want) in decoded.iter().zip(parts) {
4579 assert_eq!(got.topic, want.topic);
4580 assert_eq!(got.partition, want.partition);
4581 assert_eq!(got.error_code, want.error_code);
4582 assert_eq!(got.base_offset, want.base_offset);
4583 assert_eq!(got.log_append_time_ms, want.log_append_time_ms);
4584 if version < 5 {
4585 assert_eq!(
4586 got.log_start_offset,
4587 ProducePartitionResponse::INVALID_OFFSET,
4588 "Produce v{version} omits LogStartOffset; decode fills INVALID_OFFSET"
4589 );
4590 } else {
4591 assert_eq!(got.log_start_offset, want.log_start_offset);
4592 }
4593 assert_eq!(got.current_leader_id, MetadataResponse::NO_LEADER_ID);
4594 assert_eq!(
4595 got.current_leader_epoch,
4596 RecordBatch::NO_PARTITION_LEADER_EPOCH
4597 );
4598 if version < 8 {
4599 assert!(
4600 got.record_errors.is_empty(),
4601 "Produce v{version} omits RecordErrors; decode fills empty"
4602 );
4603 assert!(
4604 got.error_message.is_none(),
4605 "Produce v{version} omits ErrorMessage; decode fills null"
4606 );
4607 } else {
4608 assert_eq!(got.record_errors, want.record_errors);
4609 assert_eq!(got.error_message, want.error_message);
4610 }
4611 }
4612 }
4613 leftover_empty(
4614 &cur,
4615 match (version, parts.is_empty()) {
4616 (3, false) => "Produce v3 PartitionResponse.currentLeader leftover-empty",
4617 (3, true) => "Produce v3 PartitionResponse.currentLeader empty leftover-empty",
4618 (5, false) => "Produce v5 PartitionResponse.currentLeader leftover-empty",
4619 (5, true) => "Produce v5 PartitionResponse.currentLeader empty leftover-empty",
4620 (8, false) => "Produce v8 PartitionResponse.currentLeader leftover-empty",
4621 (8, true) => "Produce v8 PartitionResponse.currentLeader empty leftover-empty",
4622 (12, false) => "Produce v12 PartitionResponse.currentLeader leftover-empty",
4623 (12, true) => "Produce v12 PartitionResponse.currentLeader empty leftover-empty",
4624 _ => "Produce PartitionResponse.currentLeader leftover-empty",
4625 },
4626 )
4627 .unwrap();
4628 }
4629
4630 #[test]
4631 fn produce_v3_omitted_log_start_is_invalid_offset() {
4632 let parts = [ProducePartitionResponse {
4633 topic: "t".into(),
4634 partition: 0,
4635 error_code: 0,
4636 base_offset: 7,
4637 log_append_time_ms: RecordBatch::NO_TIMESTAMP,
4638 log_start_offset: 99,
4639 current_leader_id: -1,
4640 current_leader_epoch: -1,
4641 record_errors: Vec::new(),
4642 error_message: None,
4643 }];
4644 let mut buf = BytesMut::new();
4645 encode_produce_response(&mut buf, 3, &parts).unwrap();
4646 let mut cur = &buf[..];
4647 let (got, endpoints, ..) = decode_produce_response(&mut cur, 3).unwrap();
4648 assert!(endpoints.is_empty());
4649 assert!(cur.is_empty());
4650 let part = got.first().expect("one partition");
4651 assert_eq!(part.base_offset, 7);
4652 assert_eq!(part.log_append_time_ms, RecordBatch::NO_TIMESTAMP);
4653 assert_eq!(
4654 part.log_start_offset,
4655 ProducePartitionResponse::INVALID_OFFSET,
4656 "Produce v3 omits LogStartOffset; decode fills INVALID_OFFSET"
4657 );
4658 assert_eq!(ProducePartitionResponse::INVALID_OFFSET, -1);
4659 assert_eq!(part.current_leader_id, MetadataResponse::NO_LEADER_ID);
4660 assert_eq!(
4661 part.current_leader_epoch,
4662 RecordBatch::NO_PARTITION_LEADER_EPOCH
4663 );
4664 assert!(
4665 part.record_errors.is_empty(),
4666 "Produce v3 omits RecordErrors; decode fills empty"
4667 );
4668 assert!(
4669 part.error_message.is_none(),
4670 "Produce v3 omits ErrorMessage; decode fills null"
4671 );
4672 }
4673
4674 #[test]
4675 fn produce_record_error_matches_java() {
4676 assert_eq!(
4677 ProduceRecordError::new(1, None).to_string(),
4678 "RecordError(batchIndex=1, message=null)"
4679 );
4680 assert_eq!(
4681 ProduceRecordError::new(0, Some(String::new())).to_string(),
4682 "RecordError(batchIndex=0, message='')"
4683 );
4684 let quoted = ProduceRecordError::new(4, Some("oops".into()));
4685 assert_eq!(quoted.batch_index(), 4);
4686 assert_eq!(quoted.message(), Some("oops"));
4687 assert_eq!(
4688 quoted.to_string(),
4689 "RecordError(batchIndex=4, message='oops')"
4690 );
4691
4692 let empty = ProducePartitionResponse::partition_response("t", 0, 8);
4693 assert!(empty.record_errors.is_empty());
4694 assert!(empty.error_message.is_none());
4695 for version in [8, 9] {
4696 let mut buf = BytesMut::new();
4697 encode_produce_response(&mut buf, version, std::slice::from_ref(&empty)).unwrap();
4698 let mut cur = buf.as_ref();
4699 let (decoded, endpoints, ..) = decode_produce_response(&mut cur, version).unwrap();
4700 assert!(endpoints.is_empty());
4701 assert_eq!(decoded, std::slice::from_ref(&empty));
4702 leftover_empty(
4703 &cur,
4704 if version == 8 {
4705 "Produce v8 RecordErrors empty leftover-empty"
4706 } else {
4707 "Produce v9 RecordErrors empty leftover-empty"
4708 },
4709 )
4710 .unwrap();
4711 }
4712
4713 let mut with_errors =
4714 ProducePartitionResponse::partition_response("t", 1, crate::error::INVALID_RECORD);
4715 with_errors.record_errors = vec![
4716 ProduceRecordError::new(0, None),
4717 ProduceRecordError::new(0, Some("dup".into())),
4718 ProduceRecordError::new(2, Some("bad".into())),
4719 ];
4720 with_errors.error_message = Some("batch dropped".into());
4721 for version in [8, 9] {
4722 let mut buf = BytesMut::new();
4723 encode_produce_response(&mut buf, version, std::slice::from_ref(&with_errors)).unwrap();
4724 let mut cur = buf.as_ref();
4725 let (decoded, endpoints, ..) = decode_produce_response(&mut cur, version).unwrap();
4726 assert!(endpoints.is_empty());
4727 assert_eq!(decoded, std::slice::from_ref(&with_errors));
4728 leftover_empty(
4729 &cur,
4730 if version == 8 {
4731 "Produce v8 RecordErrors leftover-empty"
4732 } else {
4733 "Produce v9 RecordErrors leftover-empty"
4734 },
4735 )
4736 .unwrap();
4737 }
4738
4739 with_errors.current_leader_id = 2;
4740 with_errors.current_leader_epoch = 7;
4741 for version in [10, 12] {
4742 let mut buf = BytesMut::new();
4743 encode_produce_response(&mut buf, version, std::slice::from_ref(&with_errors)).unwrap();
4744 let mut cur = buf.as_ref();
4745 let (decoded, endpoints, ..) = decode_produce_response(&mut cur, version).unwrap();
4746 assert!(endpoints.is_empty());
4747 assert_eq!(decoded, std::slice::from_ref(&with_errors));
4748 leftover_empty(
4749 &cur,
4750 if version == 10 {
4751 "Produce v10 RecordErrors CurrentLeader leftover-empty"
4752 } else {
4753 "Produce v12 RecordErrors CurrentLeader leftover-empty"
4754 },
4755 )
4756 .unwrap();
4757 }
4758
4759 let mut buf = BytesMut::new();
4760 encode_produce_response(&mut buf, 3, std::slice::from_ref(&with_errors)).unwrap();
4761 let mut cur = buf.as_ref();
4762 let (decoded, ..) = decode_produce_response(&mut cur, 3).unwrap();
4763 leftover_empty(&cur, "Produce v3 RecordErrors leftover-empty").unwrap();
4764 let part = decoded.first().expect("one partition");
4765 assert!(
4766 part.record_errors.is_empty(),
4767 "Produce v3 omits RecordErrors even when the body has them"
4768 );
4769 assert!(part.error_message.is_none());
4770 assert_eq!(part.current_leader_id, MetadataResponse::NO_LEADER_ID);
4771 assert_eq!(
4772 part.current_leader_epoch,
4773 RecordBatch::NO_PARTITION_LEADER_EPOCH
4774 );
4775 }
4776
4777 #[test]
4778 fn api_versions_v3_roundtrip() {
4779 let resp = ApiVersionsResponse {
4780 error_code: 0,
4781 api_keys: vec![
4782 ApiVersion {
4783 api_key: 0,
4784 min_version: 3,
4785 max_version: 9,
4786 },
4787 ApiVersion {
4788 api_key: 3,
4789 min_version: 0,
4790 max_version: 12,
4791 },
4792 ApiVersion {
4793 api_key: 18,
4794 min_version: 0,
4795 max_version: 4,
4796 },
4797 ],
4798 throttle_time_ms: 0,
4799 ..Default::default()
4800 };
4801 let mut buf = BytesMut::new();
4802 encode_api_versions_response(&mut buf, 3, &resp).unwrap();
4803 let mut cur = &buf[..];
4804 let decoded = decode_api_versions_response(&mut cur, 3).unwrap();
4805 assert_eq!(decoded, resp);
4806 assert!(
4807 !cur.has_remaining(),
4808 "ApiVersions v3 empty features must be leftover-empty"
4809 );
4810 }
4811
4812 #[test]
4813 fn api_versions_v3_empty_features_is_zero_tagged_fields() {
4814 let resp = ApiVersionsResponse {
4815 error_code: 0,
4816 api_keys: Vec::new(),
4817 throttle_time_ms: 0,
4818 ..Default::default()
4819 };
4820 let mut buf = BytesMut::new();
4821 encode_api_versions_response(&mut buf, 3, &resp).unwrap();
4822 assert_eq!(&buf[..], &[0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00]);
4823 let mut cur = &buf[..];
4824 assert_eq!(decode_api_versions_response(&mut cur, 3).unwrap(), resp);
4825 assert!(!cur.has_remaining());
4826 }
4827
4828 #[test]
4829 fn api_versions_v3_features_roundtrip_is_leftover_empty() {
4830 const BODY: &[u8] = &[
4833 0x00, 0x00, 0x02, 0x00, 0x12, 0x00, 0x00, 0x00, 0x04, 0x00, 0x00, 0x00, 0x00, 0x00,
4834 0x03, 0x00, 0x17, 0x02, 0x11, 0x6d, 0x65, 0x74, 0x61, 0x64, 0x61, 0x74, 0x61, 0x2e,
4835 0x76, 0x65, 0x72, 0x73, 0x69, 0x6f, 0x6e, 0x00, 0x01, 0x00, 0x14, 0x00, 0x01, 0x08,
4836 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x01, 0x02, 0x17, 0x02, 0x11, 0x6d, 0x65,
4837 0x74, 0x61, 0x64, 0x61, 0x74, 0x61, 0x2e, 0x76, 0x65, 0x72, 0x73, 0x69, 0x6f, 0x6e,
4838 0x00, 0x14, 0x00, 0x01, 0x00,
4839 ];
4840 let resp = ApiVersionsResponse {
4841 error_code: 0,
4842 api_keys: vec![ApiVersion {
4843 api_key: 18,
4844 min_version: 0,
4845 max_version: 4,
4846 }],
4847 throttle_time_ms: 0,
4848 supported_features: vec![SupportedFeatureKey {
4849 name: "metadata.version".into(),
4850 min_version: 1,
4851 max_version: 20,
4852 }],
4853 finalized_features_epoch: Some(1),
4854 finalized_features: vec![FinalizedFeatureKey {
4855 name: "metadata.version".into(),
4856 max_version_level: 20,
4857 min_version_level: 1,
4858 }],
4859 zk_migration_ready: false,
4860 };
4861 let mut buf = BytesMut::new();
4862 encode_api_versions_response(&mut buf, 3, &resp).unwrap();
4863 assert_eq!(&buf[..], BODY);
4864 let mut cur = &buf[..];
4865 let decoded = decode_api_versions_response(&mut cur, 3).unwrap();
4866 assert_eq!(decoded, resp);
4867 assert!(
4868 !cur.has_remaining(),
4869 "ApiVersions v3 features must be leftover-empty"
4870 );
4871 }
4872
4873 #[test]
4874 fn api_versions_response_helpers_match_java() {
4875 assert_eq!(ApiVersionsResponse::UNKNOWN_FINALIZED_FEATURES_EPOCH, -1);
4876 assert!(!ApiVersionsResponse::should_client_throttle(1));
4877 assert!(ApiVersionsResponse::should_client_throttle(2));
4878 let resp = ApiVersionsResponse {
4879 api_keys: vec![ApiVersion {
4880 api_key: 18,
4881 min_version: 0,
4882 max_version: 4,
4883 }],
4884 zk_migration_ready: true,
4885 ..Default::default()
4886 };
4887 assert_eq!(resp.api_version(18).map(|v| v.max_version), Some(4));
4888 assert!(resp.api_version(1).is_none());
4889 assert!(resp.zk_migration_ready());
4890 let produce = ApiVersion {
4891 api_key: 0,
4892 min_version: 0,
4893 max_version: 12,
4894 };
4895 assert_eq!(produce.api_key(), 0);
4896 assert_eq!(produce.min_version(), 0);
4897 assert_eq!(produce.max_version(), 12);
4898 let overlap = ApiVersion {
4899 api_key: 0,
4900 min_version: 3,
4901 max_version: 9,
4902 };
4903 let got = ApiVersionsResponse::intersect(Some(&produce), Some(&overlap))
4904 .expect("same api key")
4905 .expect("overlap");
4906 assert_eq!(got.api_key(), 0);
4907 assert_eq!(got.min_version(), 3);
4908 assert_eq!(got.max_version(), 9);
4909 let disjoint = ApiVersion {
4910 api_key: 0,
4911 min_version: 13,
4912 max_version: 15,
4913 };
4914 assert_eq!(
4915 ApiVersionsResponse::intersect(Some(&produce), Some(&disjoint)).expect("same api key"),
4916 None
4917 );
4918 assert_eq!(
4919 ApiVersionsResponse::intersect(None, Some(&produce)).expect("null"),
4920 None
4921 );
4922 assert_eq!(
4923 ApiVersionsResponse::intersect(Some(&produce), None).expect("null"),
4924 None
4925 );
4926 assert_eq!(
4927 ApiVersionsResponse::intersect(None, None).expect("null"),
4928 None
4929 );
4930 let fetch = ApiVersion {
4931 api_key: 1,
4932 min_version: 0,
4933 max_version: 17,
4934 };
4935 let err = ApiVersionsResponse::intersect(Some(&produce), Some(&fetch)).unwrap_err();
4936 assert!(err
4937 .to_string()
4938 .contains("thisVersion.apiKey: 0 must be equal to other.apiKey: 1"));
4939 let supported = SupportedFeatureKey {
4940 name: "metadata.version".into(),
4941 min_version: 1,
4942 max_version: 20,
4943 };
4944 assert_eq!(supported.name(), "metadata.version");
4945 assert_eq!(supported.min_version(), 1);
4946 assert_eq!(supported.max_version(), 20);
4947 let finalized = FinalizedFeatureKey {
4948 name: "metadata.version".into(),
4949 max_version_level: 20,
4950 min_version_level: 1,
4951 };
4952 assert_eq!(finalized.name(), "metadata.version");
4953 assert_eq!(finalized.max_version_level(), 20);
4954 assert_eq!(finalized.min_version_level(), 1);
4955 }
4956
4957 #[test]
4958 fn api_versions_response_error_counts_matches_java() {
4959 let none = ApiVersionsResponse {
4967 error_code: 0,
4968 api_keys: vec![ApiVersion {
4969 api_key: 18,
4970 min_version: 0,
4971 max_version: 4,
4972 }],
4973 throttle_time_ms: 0,
4974 supported_features: vec![SupportedFeatureKey {
4975 name: "metadata.version".into(),
4976 min_version: 1,
4977 max_version: 20,
4978 }],
4979 finalized_features_epoch: Some(1),
4980 finalized_features: vec![FinalizedFeatureKey {
4981 name: "metadata.version".into(),
4982 max_version_level: 20,
4983 min_version_level: 1,
4984 }],
4985 zk_migration_ready: false,
4986 };
4987 assert_eq!(
4988 none.error_counts(),
4989 HashMap::from([(0, 1)]),
4990 "NONE is a singleton 1, not an empty map"
4991 );
4992 let full = ApiVersionsResponse {
4993 error_code: crate::error::INVALID_REQUEST,
4994 api_keys: vec![ApiVersion {
4995 api_key: 18,
4996 min_version: 0,
4997 max_version: 4,
4998 }],
4999 throttle_time_ms: 0,
5000 supported_features: vec![SupportedFeatureKey {
5001 name: "metadata.version".into(),
5002 min_version: 1,
5003 max_version: 20,
5004 }],
5005 finalized_features_epoch: Some(1),
5006 finalized_features: vec![FinalizedFeatureKey {
5007 name: "metadata.version".into(),
5008 max_version_level: 20,
5009 min_version_level: 1,
5010 }],
5011 zk_migration_ready: false,
5012 };
5013 assert_eq!(
5014 full.error_counts(),
5015 HashMap::from([(crate::error::INVALID_REQUEST, 1)])
5016 );
5017 for version in 0..=4_i16 {
5018 let mut resp = BytesMut::new();
5019 encode_api_versions_response(&mut resp, version, &full).unwrap();
5020 let mut cur = &resp[..];
5021 let decoded = decode_api_versions_response(&mut cur, version).unwrap();
5022 assert_eq!(
5023 decoded.error_counts(),
5024 HashMap::from([(crate::error::INVALID_REQUEST, 1)]),
5025 "ApiVersions v{version} errorCounts must count the decoded code"
5026 );
5027 assert!(
5028 cur.is_empty(),
5029 "ApiVersions v{version} errorCounts leftover-empty; leftover {} bytes",
5030 cur.len()
5031 );
5032 }
5033 }
5034
5035 #[test]
5036 fn api_versions_request_is_valid_matches_java() {
5037 assert!(ApiVersionsRequest::is_valid(2, "", ""));
5038 assert!(ApiVersionsRequest::is_valid(2, "-invalid", "x."));
5039 assert!(ApiVersionsRequest::is_valid(
5040 3,
5041 crate::CLIENT_NAME,
5042 crate::CLIENT_VERSION
5043 ));
5044 assert!(ApiVersionsRequest::is_valid(4, "a", "1"));
5045 assert!(ApiVersionsRequest::is_valid(3, "a-b.c", "0.1.0"));
5046 assert!(!ApiVersionsRequest::is_valid(3, "", "0.1.0"));
5047 assert!(!ApiVersionsRequest::is_valid(3, "partitionline", ""));
5048 assert!(!ApiVersionsRequest::is_valid(3, "-x", "0.1.0"));
5049 assert!(!ApiVersionsRequest::is_valid(3, "x-", "0.1.0"));
5050 assert!(!ApiVersionsRequest::is_valid(3, ".x", "0.1.0"));
5051 assert!(!ApiVersionsRequest::is_valid(3, "x.", "0.1.0"));
5052 assert!(!ApiVersionsRequest::is_valid(3, "x_y", "0.1.0"));
5053 }
5054
5055 #[test]
5056 fn api_versions_request_error_response_matches_java() {
5057 let none = ApiVersionsRequest::error_response(0);
5062 assert_eq!(none.error_code, 0);
5063 assert!(none.api_keys.is_empty());
5064 assert_eq!(none.throttle_time_ms, 0);
5065 assert!(none.supported_features.is_empty());
5066 assert!(none.finalized_features.is_empty());
5067 assert!(none.finalized_features_epoch.is_none());
5068 assert!(!none.zk_migration_ready);
5069
5070 let other = ApiVersionsRequest::error_response(crate::error::INVALID_REQUEST);
5071 assert_eq!(other.error_code, crate::error::INVALID_REQUEST);
5072 assert!(other.api_keys.is_empty());
5073
5074 let unsupported = ApiVersionsRequest::error_response(crate::error::UNSUPPORTED_VERSION);
5075 assert_eq!(unsupported.error_code, crate::error::UNSUPPORTED_VERSION);
5076 assert_eq!(
5077 unsupported.api_keys,
5078 [ApiVersion {
5079 api_key: API_VERSIONS,
5080 min_version: 0,
5081 max_version: 4,
5082 }]
5083 );
5084 assert_eq!(
5085 unsupported.api_version(API_VERSIONS).map(|v| v.max_version),
5086 Some(4)
5087 );
5088
5089 for version in [0_i16, 1, 3, 4] {
5090 let mut buf = BytesMut::new();
5091 encode_api_versions_response(&mut buf, version, &unsupported).unwrap();
5092 let mut cur = buf.as_ref();
5093 let decoded = decode_api_versions_response(&mut cur, version).unwrap();
5094 assert_eq!(decoded.error_code, crate::error::UNSUPPORTED_VERSION);
5095 assert_eq!(decoded.api_keys, unsupported.api_keys);
5096 leftover_empty(
5097 &cur,
5098 match version {
5099 0 => "ApiVersions v0 Request.getErrorResponse leftover-empty",
5100 1 => "ApiVersions v1 Request.getErrorResponse leftover-empty",
5101 3 => "ApiVersions v3 Request.getErrorResponse leftover-empty",
5102 _ => "ApiVersions v4 Request.getErrorResponse leftover-empty",
5103 },
5104 )
5105 .unwrap();
5106 }
5107 for version in [0_i16, 1, 3, 4] {
5108 let mut buf = BytesMut::new();
5109 encode_api_versions_response(&mut buf, version, &none).unwrap();
5110 let mut cur = buf.as_ref();
5111 let decoded = decode_api_versions_response(&mut cur, version).unwrap();
5112 assert_eq!(decoded.error_code, 0);
5113 assert!(decoded.api_keys.is_empty());
5114 leftover_empty(
5115 &cur,
5116 match version {
5117 0 => "ApiVersions v0 Request.getErrorResponse empty leftover-empty",
5118 1 => "ApiVersions v1 Request.getErrorResponse empty leftover-empty",
5119 3 => "ApiVersions v3 Request.getErrorResponse empty leftover-empty",
5120 _ => "ApiVersions v4 Request.getErrorResponse empty leftover-empty",
5121 },
5122 )
5123 .unwrap();
5124 }
5125 }
5126
5127 #[test]
5128 fn api_versions_create_finalized_feature_keys_matches_java() {
5129 let empty =
5133 ApiVersionsResponse::create_finalized_feature_keys(std::iter::empty::<(&str, i16)>());
5134 assert!(empty.is_empty());
5135 assert!(ApiVersionsResponse::create_finalized_feature_keys([("f", 0)]).is_empty());
5136
5137 let skipped = ApiVersionsResponse::create_finalized_feature_keys([("f", 5), ("f", 0)]);
5138 assert!(skipped.is_empty(), "last-wins then skip");
5139
5140 let kept = ApiVersionsResponse::create_finalized_feature_keys([("f", 0), ("f", 5)]);
5141 assert_eq!(
5142 kept,
5143 [FinalizedFeatureKey {
5144 name: "f".into(),
5145 max_version_level: 5,
5146 min_version_level: 5,
5147 }]
5148 );
5149
5150 let keys = ApiVersionsResponse::create_finalized_feature_keys([
5151 ("b", 2),
5152 ("a", 0),
5153 ("c", 3),
5154 ("b", 7),
5155 ]);
5156 assert_eq!(
5157 keys,
5158 [
5159 FinalizedFeatureKey {
5160 name: "b".into(),
5161 max_version_level: 7,
5162 min_version_level: 7,
5163 },
5164 FinalizedFeatureKey {
5165 name: "c".into(),
5166 max_version_level: 3,
5167 min_version_level: 3,
5168 },
5169 ],
5170 "first-seen order; skip 0"
5171 );
5172
5173 for version in [3_i16, 4] {
5174 let resp = ApiVersionsResponse {
5175 finalized_features: keys.clone(),
5176 ..Default::default()
5177 };
5178 let mut buf = BytesMut::new();
5179 encode_api_versions_response(&mut buf, version, &resp).unwrap();
5180 let mut cur = buf.as_ref();
5181 let decoded = decode_api_versions_response(&mut cur, version).unwrap();
5182 assert_eq!(decoded.finalized_features, keys);
5183 leftover_empty(
5184 &cur,
5185 match version {
5186 3 => "ApiVersions v3 createFinalizedFeatureKeys leftover-empty",
5187 _ => "ApiVersions v4 createFinalizedFeatureKeys leftover-empty",
5188 },
5189 )
5190 .unwrap();
5191 }
5192 for version in [3_i16, 4] {
5193 let resp = ApiVersionsResponse {
5194 finalized_features: empty.clone(),
5195 ..Default::default()
5196 };
5197 let mut buf = BytesMut::new();
5198 encode_api_versions_response(&mut buf, version, &resp).unwrap();
5199 let mut cur = buf.as_ref();
5200 let decoded = decode_api_versions_response(&mut cur, version).unwrap();
5201 assert!(decoded.finalized_features.is_empty());
5202 leftover_empty(
5203 &cur,
5204 match version {
5205 3 => "ApiVersions v3 createFinalizedFeatureKeys empty leftover-empty",
5206 _ => "ApiVersions v4 createFinalizedFeatureKeys empty leftover-empty",
5207 },
5208 )
5209 .unwrap();
5210 }
5211 }
5212
5213 #[test]
5214 fn api_versions_maybe_filter_supported_feature_keys_matches_java() {
5215 let kraft = SupportedFeatureKey {
5220 name: "kraft.version".into(),
5221 min_version: 0,
5222 max_version: 1,
5223 };
5224 let meta = SupportedFeatureKey {
5225 name: "metadata.version".into(),
5226 min_version: 1,
5227 max_version: 20,
5228 };
5229 assert!(ApiVersionsResponse::maybe_filter_supported_feature_keys(&[], true).is_empty());
5230 assert!(ApiVersionsResponse::maybe_filter_supported_feature_keys(&[], false).is_empty());
5231
5232 let filtered = ApiVersionsResponse::maybe_filter_supported_feature_keys(
5233 &[meta.clone(), kraft.clone()],
5234 true,
5235 );
5236 assert_eq!(filtered.as_slice(), std::slice::from_ref(&meta));
5237
5238 let kept = ApiVersionsResponse::maybe_filter_supported_feature_keys(
5239 &[meta.clone(), kraft.clone()],
5240 false,
5241 );
5242 assert_eq!(kept, [meta.clone(), kraft.clone()]);
5243
5244 let skipped = ApiVersionsResponse::maybe_filter_supported_feature_keys(
5245 std::slice::from_ref(&kraft),
5246 true,
5247 );
5248 assert!(skipped.is_empty());
5249
5250 let dup = ApiVersionsResponse::maybe_filter_supported_feature_keys(
5251 &[meta.clone(), meta.clone()],
5252 true,
5253 );
5254 assert_eq!(dup.len(), 2, "names are not uniqued");
5255
5256 for version in [3_i16, 4] {
5257 let resp = ApiVersionsResponse {
5258 supported_features: filtered.clone(),
5259 ..Default::default()
5260 };
5261 let mut buf = BytesMut::new();
5262 encode_api_versions_response(&mut buf, version, &resp).unwrap();
5263 let mut cur = buf.as_ref();
5264 let decoded = decode_api_versions_response(&mut cur, version).unwrap();
5265 assert_eq!(decoded.supported_features, filtered);
5266 leftover_empty(
5267 &cur,
5268 match version {
5269 3 => "ApiVersions v3 maybeFilterSupportedFeatureKeys leftover-empty",
5270 _ => "ApiVersions v4 maybeFilterSupportedFeatureKeys leftover-empty",
5271 },
5272 )
5273 .unwrap();
5274 }
5275 for version in [3_i16, 4] {
5276 let resp = ApiVersionsResponse::default();
5277 let mut buf = BytesMut::new();
5278 encode_api_versions_response(&mut buf, version, &resp).unwrap();
5279 let mut cur = buf.as_ref();
5280 let decoded = decode_api_versions_response(&mut cur, version).unwrap();
5281 assert!(decoded.supported_features.is_empty());
5282 leftover_empty(
5283 &cur,
5284 match version {
5285 3 => "ApiVersions v3 maybeFilterSupportedFeatureKeys empty leftover-empty",
5286 _ => "ApiVersions v4 maybeFilterSupportedFeatureKeys empty leftover-empty",
5287 },
5288 )
5289 .unwrap();
5290 }
5291 }
5292
5293 #[test]
5294 fn api_versions_v3_matches_v4_and_does_not_speak_v5() {
5295 let mut v3 = BytesMut::new();
5300 encode_api_versions_request(&mut v3, 3, "partitionline", "0.1.0").unwrap();
5301 let mut v4 = BytesMut::new();
5302 encode_api_versions_request(&mut v4, 4, "partitionline", "0.1.0").unwrap();
5303 assert_eq!(v3.as_ref(), v4.as_ref(), "v3 and v4 request bodies match");
5304 let mut v0 = BytesMut::new();
5305 encode_api_versions_request(&mut v0, 0, "partitionline", "0.1.0").unwrap();
5306 assert!(v0.is_empty(), "v0–v2 request is empty");
5307 encode_api_versions_request(&mut v0, 2, "partitionline", "0.1.0").unwrap();
5308 assert!(v0.is_empty(), "v2 request is empty");
5309 let mut empty: &[u8] = &[];
5310 assert_eq!(
5311 decode_api_versions_request(&mut empty, 0).unwrap(),
5312 (String::new(), String::new())
5313 );
5314 let mut empty: &[u8] = &[];
5315 assert_eq!(
5316 decode_api_versions_request(&mut empty, 2).unwrap(),
5317 (String::new(), String::new())
5318 );
5319 let mut cur = v3.as_ref();
5320 assert_eq!(
5321 decode_api_versions_request(&mut cur, 3).unwrap(),
5322 ("partitionline".into(), "0.1.0".into())
5323 );
5324 assert!(!cur.has_remaining(), "v3 request leftover-empty");
5325 let mut cur = v4.as_ref();
5326 assert_eq!(
5327 decode_api_versions_request(&mut cur, 4).unwrap(),
5328 ("partitionline".into(), "0.1.0".into())
5329 );
5330 assert!(!cur.has_remaining(), "v4 request leftover-empty");
5331 let err = encode_api_versions_request(&mut BytesMut::new(), 5, "partitionline", "0.1.0")
5332 .unwrap_err();
5333 assert!(
5334 err.to_string().contains("not implemented"),
5335 "v5 is not spoken, got {err}"
5336 );
5337 let mut empty: &[u8] = &[];
5338 let err = decode_api_versions_request(&mut empty, 5).unwrap_err();
5339 assert!(
5340 err.to_string().contains("not implemented"),
5341 "v5 decode is not spoken, got {err}"
5342 );
5343 assert_eq!(crate::protocol::api_keys::pick_version(0, 3, 0, 4), Some(3));
5344 assert_eq!(crate::protocol::api_keys::pick_version(0, 4, 0, 4), Some(4));
5345 assert_eq!(crate::protocol::api_keys::pick_version(5, 5, 0, 4), None);
5346
5347 let kraft = SupportedFeatureKey {
5348 name: "kraft.version".into(),
5349 min_version: 0,
5350 max_version: 1,
5351 };
5352 let meta = SupportedFeatureKey {
5353 name: "metadata.version".into(),
5354 min_version: 1,
5355 max_version: 20,
5356 };
5357 let resp = ApiVersionsResponse {
5358 error_code: 0,
5359 api_keys: vec![ApiVersion {
5360 api_key: 18,
5361 min_version: 0,
5362 max_version: 4,
5363 }],
5364 throttle_time_ms: 0,
5365 supported_features: vec![meta.clone(), kraft.clone()],
5366 finalized_features_epoch: None,
5367 finalized_features: Vec::new(),
5368 zk_migration_ready: false,
5369 };
5370 v3.clear();
5371 encode_api_versions_response(&mut v3, 3, &resp).unwrap();
5372 v4.clear();
5373 encode_api_versions_response(&mut v4, 4, &resp).unwrap();
5374 assert_ne!(
5375 v3.as_ref(),
5376 v4.as_ref(),
5377 "v3 omits SupportedFeatures with MinVersion 0"
5378 );
5379 let mut cur = v3.as_ref();
5380 let decoded = decode_api_versions_response(&mut cur, 3).unwrap();
5381 assert_eq!(decoded.supported_features, vec![meta.clone()]);
5382 assert!(!cur.has_remaining(), "v3 response leftover-empty");
5383 let mut cur = v4.as_ref();
5384 let decoded = decode_api_versions_response(&mut cur, 4).unwrap();
5385 assert_eq!(decoded.supported_features, vec![meta, kraft]);
5386 assert!(!cur.has_remaining(), "v4 response leftover-empty");
5387 v3.clear();
5388 let err = encode_api_versions_response(&mut v3, 5, &resp).unwrap_err();
5389 assert!(
5390 err.to_string().contains("not implemented"),
5391 "v5 response is not spoken, got {err}"
5392 );
5393 }
5394
5395 #[test]
5396 fn api_versions_kip511_v0_unsupported_body_parses_when_sent_v4() {
5397 let resp = ApiVersionsResponse {
5401 error_code: crate::error::UNSUPPORTED_VERSION,
5402 api_keys: vec![ApiVersion {
5403 api_key: API_VERSIONS,
5404 min_version: 0,
5405 max_version: 3,
5406 }],
5407 ..Default::default()
5408 };
5409 let mut buf = BytesMut::new();
5410 encode_api_versions_response(&mut buf, 0, &resp).unwrap();
5411 let decoded = decode_api_versions_handshake(&buf, 4).unwrap();
5412 assert_eq!(decoded.error_code, crate::error::UNSUPPORTED_VERSION);
5413 assert_eq!(decoded.api_keys, resp.api_keys);
5414 assert_eq!(crate::protocol::api_keys::pick_version(0, 3, 0, 4), Some(3));
5415 assert_eq!(crate::protocol::api_keys::pick_version(0, 0, 0, 4), Some(0));
5416 }
5417
5418 #[test]
5419 fn produce_v3_transactional_id_is_not_null() {
5420 let rec = Record {
5421 offset: 0,
5422 timestamp: 1,
5423 key: None,
5424 value: Some(Bytes::from_static(b"x")),
5425 headers: vec![],
5426 };
5427 let topics = vec![ProduceTopicData {
5428 topic: "t".into(),
5429 partitions: vec![ProducePartitionData {
5430 index: 0,
5431 records: RecordBatch::from_records(vec![rec]),
5432 }],
5433 }];
5434 let mut buf = BytesMut::new();
5435 encode_produce_request(&mut buf, 3, Some("tx-1"), -1, 1000, &topics).unwrap();
5436 let (txn, acks, _, _) = decode_produce_request(&mut &buf[..], 3).unwrap();
5437 assert_eq!(txn.as_deref(), Some("tx-1"));
5438 assert_eq!(acks, -1);
5439 }
5440
5441 #[test]
5442 fn produce_v9_roundtrip_is_leftover_empty() {
5443 let rec = Record {
5444 offset: 0,
5445 timestamp: 42,
5446 key: None,
5447 value: Some(Bytes::from_static(b"hi")),
5448 headers: vec![],
5449 };
5450 let topics = vec![ProduceTopicData {
5451 topic: "t".into(),
5452 partitions: vec![ProducePartitionData {
5453 index: 0,
5454 records: RecordBatch::from_records(vec![rec]),
5455 }],
5456 }];
5457 let mut buf = BytesMut::new();
5458 encode_produce_request(&mut buf, 9, None, 1, 1500, &topics).unwrap();
5459 let mut cur = &buf[..];
5460 let (txn, acks, timeout, decoded) = decode_produce_request(&mut cur, 9).unwrap();
5461 assert_eq!(txn, None);
5462 assert_eq!(acks, 1);
5463 assert_eq!(timeout, 1500);
5464 assert_eq!(decoded[0].topic, "t");
5465 assert_eq!(
5466 decoded[0].partitions[0].records.records[0].value.as_deref(),
5467 Some(&b"hi"[..])
5468 );
5469 assert!(
5470 cur.is_empty(),
5471 "Produce v9 request must consume compact tagged fields"
5472 );
5473
5474 buf.clear();
5475 encode_produce_request(&mut buf, 12, None, 1, 1500, &topics).unwrap();
5476 let mut cur = &buf[..];
5477 let (txn, acks, timeout, decoded) = decode_produce_request(&mut cur, 12).unwrap();
5478 assert_eq!(txn, None);
5479 assert_eq!(acks, 1);
5480 assert_eq!(timeout, 1500);
5481 assert_eq!(decoded[0].topic, "t");
5482 assert!(
5483 cur.is_empty(),
5484 "Produce v12 request must consume compact tagged fields"
5485 );
5486
5487 let mut txn_buf = BytesMut::new();
5488 encode_produce_request(&mut txn_buf, 8, Some("tx-1"), 1, 1500, &topics).unwrap();
5489 let mut cur = &txn_buf[..];
5490 let (txn, _, _, _) = decode_produce_request(&mut cur, 8).unwrap();
5491 assert_eq!(txn.as_deref(), Some("tx-1"));
5492 assert!(
5493 cur.is_empty(),
5494 "Produce v8 request leftover {} bytes",
5495 cur.len()
5496 );
5497
5498 buf.clear();
5499 assert!(
5500 encode_produce_request(&mut buf, 13, None, 1, 1500, &topics).is_err(),
5501 "Produce v13+ (topic IDs) is not spoken"
5502 );
5503 }
5504
5505 #[test]
5506 fn produce_v9_response_matches_compact_layout() {
5507 const RESP: &[u8] = &[
5511 0x02, 0x02, 0x74, 0x02, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
5512 0x00, 0x00, 0x00, 0x00, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0x00, 0x00,
5513 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
5514 0x00,
5515 ];
5516 let parts = [ProducePartitionResponse {
5517 topic: "t".into(),
5518 partition: 0,
5519 error_code: 0,
5520 base_offset: 0,
5521 log_append_time_ms: RecordBatch::NO_TIMESTAMP,
5522 log_start_offset: 0,
5523 current_leader_id: -1,
5524 current_leader_epoch: -1,
5525 record_errors: Vec::new(),
5526 error_message: None,
5527 }];
5528 let mut buf = BytesMut::new();
5529 encode_produce_response(&mut buf, 9, &parts).unwrap();
5530 assert_eq!(&buf[..], RESP);
5531 let mut cur = &buf[..];
5532 let got = decode_produce_response(&mut cur, 9).unwrap();
5533 assert_eq!(got.0, parts);
5534 assert!(got.1.is_empty());
5535 assert!(
5536 cur.is_empty(),
5537 "Produce v9 response must consume compact tagged fields"
5538 );
5539
5540 buf.clear();
5541 encode_produce_response(&mut buf, 12, &parts).unwrap();
5542 assert_eq!(
5543 &buf[..],
5544 RESP,
5545 "Produce v12 with omitted CurrentLeader matches v9 bytes"
5546 );
5547 let mut cur = &buf[..];
5548 let got = decode_produce_response(&mut cur, 12).unwrap();
5549 assert_eq!(got.0, parts);
5550 assert!(got.1.is_empty());
5551 assert!(
5552 cur.is_empty(),
5553 "Produce v12 empty CurrentLeader must consume compact tagged fields"
5554 );
5555 }
5556
5557 #[test]
5558 fn produce_v11_current_leader_tagged_is_leftover_empty() {
5559 const RESP: &[u8] = &[
5562 0x02, 0x02, 0x74, 0x02, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
5563 0x00, 0x00, 0x00, 0x00, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0x00, 0x00,
5564 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x01, 0x00, 0x09, 0x00, 0x00, 0x00,
5565 0x02, 0x00, 0x00, 0x00, 0x07, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
5566 ];
5567 let parts = [ProducePartitionResponse {
5568 topic: "t".into(),
5569 partition: 0,
5570 error_code: 0,
5571 base_offset: 0,
5572 log_append_time_ms: RecordBatch::NO_TIMESTAMP,
5573 log_start_offset: 0,
5574 current_leader_id: 2,
5575 current_leader_epoch: 7,
5576 record_errors: Vec::new(),
5577 error_message: None,
5578 }];
5579 let mut buf = BytesMut::new();
5580 encode_produce_response(&mut buf, 11, &parts).unwrap();
5581 assert_eq!(&buf[..], RESP);
5582 let mut cur = &buf[..];
5583 let got = decode_produce_response(&mut cur, 11).unwrap();
5584 assert_eq!(got.0, parts);
5585 assert!(got.1.is_empty());
5586 assert!(
5587 cur.is_empty(),
5588 "Produce v11 CurrentLeader must consume nested tagged fields"
5589 );
5590 buf.clear();
5591 encode_produce_response(&mut buf, 10, &parts).unwrap();
5592 assert_eq!(&buf[..], RESP, "Produce v10 CurrentLeader matches v11");
5593 buf.clear();
5594 encode_produce_response(&mut buf, 12, &parts).unwrap();
5595 assert_eq!(&buf[..], RESP, "Produce v12 CurrentLeader matches v11");
5596 buf.clear();
5597 assert!(
5598 encode_produce_response(&mut buf, 13, &parts).is_err(),
5599 "Produce v13+ (topic IDs) is not spoken"
5600 );
5601 }
5602
5603 #[test]
5604 fn produce_v10_node_endpoints_tagged_is_leftover_empty() {
5605 let parts = [ProducePartitionResponse {
5606 topic: "t".into(),
5607 partition: 0,
5608 error_code: 6,
5609 base_offset: ProducePartitionResponse::INVALID_OFFSET,
5610 log_append_time_ms: RecordBatch::NO_TIMESTAMP,
5611 log_start_offset: 0,
5612 current_leader_id: 3,
5613 current_leader_epoch: 1,
5614 record_errors: Vec::new(),
5615 error_message: None,
5616 }];
5617 let endpoints = [NodeEndpoint {
5618 node_id: 3,
5619 host: "h".into(),
5620 port: 1,
5621 rack: None,
5622 }];
5623 let mut buf = BytesMut::new();
5624 encode_produce_response_with_endpoints(&mut buf, 10, &parts, &endpoints).unwrap();
5625 let mut cur = &buf[..];
5626 let (got, eps, ..) = decode_produce_response(&mut cur, 10).unwrap();
5627 assert_eq!(got[0].current_leader_id, 3);
5628 assert_eq!(eps, endpoints);
5629 assert!(
5630 cur.is_empty(),
5631 "Produce NodeEndpoints tagged field 0 must consume nested tagged fields"
5632 );
5633 let mut omitted = BytesMut::new();
5634 encode_produce_response(&mut omitted, 10, &parts).unwrap();
5635 assert_ne!(
5636 &buf[..],
5637 &omitted[..],
5638 "NodeEndpoints tagged field 0 must not equal empty tags"
5639 );
5640 let mut v9 = BytesMut::new();
5641 encode_produce_response_with_endpoints(&mut v9, 9, &parts, &endpoints).unwrap();
5642 let mut empty = BytesMut::new();
5643 encode_produce_response(&mut empty, 9, &parts).unwrap();
5644 assert_eq!(&v9[..], &empty[..], "Produce v9 must omit NodeEndpoints");
5645 }
5646
5647 #[test]
5648 fn metadata_broker_and_node_endpoint_match_java_node() {
5649 let broker = Broker::new(1, "127.0.0.1", 9092, Some("r".into()));
5650 assert_eq!(broker.id(), 1);
5651 assert_eq!(broker.id_string(), "1");
5652 assert_eq!(broker.host(), "127.0.0.1");
5653 assert_eq!(broker.port(), 9092);
5654 assert_eq!(broker.rack(), Some("r"));
5655 assert!(broker.has_rack());
5656 assert!(!broker.is_fenced());
5657 assert!(!broker.is_empty());
5658 assert_eq!(
5659 broker.to_string(),
5660 "127.0.0.1:9092 (id: 1 rack: r isFenced: false)"
5661 );
5662 let endpoint = NodeEndpoint::from(broker.clone());
5663 assert_eq!(endpoint.id(), 1);
5664 assert_eq!(endpoint.host(), "127.0.0.1");
5665 assert_eq!(endpoint.port(), 9092);
5666 assert_eq!(endpoint.rack(), Some("r"));
5667 assert!(endpoint.has_rack());
5668 assert!(!endpoint.is_fenced());
5669 assert_eq!(endpoint.to_string(), broker.to_string());
5670 assert_eq!(Broker::from(endpoint.clone()), broker);
5671 let empty = Broker::no_node();
5672 assert_eq!(empty.id(), -1);
5673 assert!(empty.is_empty());
5674 assert_eq!(empty.id_string(), "-1");
5675 assert_eq!(empty.to_string(), ":-1 (id: -1 rack: null isFenced: false)");
5676 let empty_ep = NodeEndpoint::no_node();
5677 assert!(empty_ep.is_empty());
5678 assert_eq!(empty_ep.to_string(), empty.to_string());
5679 }
5680
5681 #[test]
5682 fn metadata_v12_roundtrip() {
5683 let resp = MetadataResponse {
5684 throttle_time_ms: 0,
5685 brokers: vec![Broker {
5686 node_id: 1,
5687 host: "127.0.0.1".into(),
5688 port: 9092,
5689 rack: None,
5690 }],
5691 cluster_id: Some("cid".into()),
5692 controller_id: 1,
5693 topics: vec![TopicMetadata {
5694 error_code: 0,
5695 name: Some("orders".into()),
5696 topic_id: [1u8; 16],
5697 is_internal: false,
5698 partitions: vec![PartitionMetadata {
5699 error_code: 0,
5700 partition_index: 0,
5701 leader_id: 1,
5702 leader_epoch: 3,
5703 replica_nodes: vec![1],
5704 isr_nodes: vec![1],
5705 offline_replicas: vec![2],
5706 }],
5707 topic_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
5708 }],
5709 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
5710 error_code: 0,
5711 };
5712 let mut buf = BytesMut::new();
5713 encode_metadata_response(&mut buf, 12, &resp).unwrap();
5714 let mut cur = &buf[..];
5715 let decoded = decode_metadata_response(&mut cur, 12).unwrap();
5716 assert_eq!(decoded, resp);
5717 leftover_empty(&cur, "Metadata v12").unwrap();
5718
5719 let mut v13 = BytesMut::new();
5720 encode_metadata_response(&mut v13, 13, &resp).unwrap();
5721 let mut cur = &v13[..];
5722 let decoded = decode_metadata_response(&mut cur, 13).unwrap();
5723 assert_eq!(decoded, resp);
5724 leftover_empty(&cur, "Metadata v13").unwrap();
5725 assert_ne!(
5726 &buf[..],
5727 &v13[..],
5728 "Metadata v13 must write top-level ErrorCode before tagged fields"
5729 );
5730 }
5731
5732 #[test]
5733 fn metadata_v13_top_error_fails_check() {
5734 let resp = MetadataResponse {
5735 throttle_time_ms: 0,
5736 brokers: Vec::new(),
5737 cluster_id: None,
5738 controller_id: MetadataResponse::NO_CONTROLLER_ID,
5739 topics: Vec::new(),
5740 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
5741 error_code: crate::error::UNKNOWN_TOPIC_OR_PARTITION,
5742 };
5743 assert_eq!(resp.controller_id, MetadataResponse::NO_CONTROLLER_ID);
5744 assert_eq!(MetadataResponse::NO_CONTROLLER_ID, -1);
5745 assert_eq!(MetadataResponse::NO_LEADER_ID, -1);
5746 let mut buf = BytesMut::new();
5747 encode_metadata_response(&mut buf, 13, &resp).unwrap();
5748 let decoded = decode_metadata_response(&mut &buf[..], 13).unwrap();
5749 assert_eq!(decoded.error_code, crate::error::UNKNOWN_TOPIC_OR_PARTITION);
5750 assert_eq!(
5751 decoded.check().unwrap_err().broker_code(),
5752 Some(crate::error::UNKNOWN_TOPIC_OR_PARTITION)
5753 );
5754 }
5755
5756 #[test]
5757 fn metadata_has_reliable_leader_epochs_matches_java() {
5758 assert!(!MetadataResponse::has_reliable_leader_epochs(8));
5759 assert!(MetadataResponse::has_reliable_leader_epochs(9));
5760 assert!(MetadataResponse::has_reliable_leader_epochs(13));
5761 assert_eq!(MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED, i32::MIN);
5762 assert_eq!(
5763 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
5764 crate::AUTHORIZED_OPERATIONS_OMITTED
5765 );
5766 assert!(!MetadataResponse::should_client_throttle(5));
5767 assert!(MetadataResponse::should_client_throttle(6));
5768 let with_epoch = PartitionMetadata {
5769 error_code: 0,
5770 partition_index: 1,
5771 leader_id: 2,
5772 leader_epoch: 8,
5773 replica_nodes: vec![2, 3],
5774 isr_nodes: vec![2],
5775 offline_replicas: vec![3],
5776 };
5777 let stripped = with_epoch.without_leader_epoch();
5778 assert_eq!(
5779 stripped.leader_epoch,
5780 RecordBatch::NO_PARTITION_LEADER_EPOCH
5781 );
5782 assert_eq!(stripped.error_code, with_epoch.error_code);
5783 assert_eq!(stripped.partition_index, with_epoch.partition_index);
5784 assert_eq!(stripped.leader_id, with_epoch.leader_id);
5785 assert_eq!(stripped.replica_nodes, with_epoch.replica_nodes);
5786 assert_eq!(stripped.isr_nodes, with_epoch.isr_nodes);
5787 assert_eq!(stripped.offline_replicas, with_epoch.offline_replicas);
5788 assert_eq!(with_epoch.leader_epoch, 8);
5789 let topic = TopicMetadata {
5790 error_code: 0,
5791 name: Some("t".into()),
5792 topic_id: [0; 16],
5793 is_internal: false,
5794 partitions: vec![with_epoch.clone()],
5795 topic_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
5796 };
5797 assert_eq!(
5798 topic.to_string(),
5799 "TopicMetadata{error=NONE, topic='t', topicId='AAAAAAAAAAAAAAAAAAAAAA', isInternal=false, partitionMetadata=[PartitionMetadata(error=NONE, partition=t-1, leader=Optional[2], leaderEpoch=Optional[8], replicas=2,3, isr=2, offlineReplicas=3)], authorizedOperations=-2147483648}"
5800 );
5801 let empty = TopicMetadata::error(0, None, [0; 16]);
5802 assert_eq!(
5803 empty.to_string(),
5804 "TopicMetadata{error=NONE, topic='', topicId='AAAAAAAAAAAAAAAAAAAAAA', isInternal=false, partitionMetadata=[], authorizedOperations=-2147483648}"
5805 );
5806 let unnamed = TopicMetadata {
5807 error_code: crate::error::UNKNOWN_TOPIC_OR_PARTITION,
5808 name: None,
5809 topic_id: [0; 16],
5810 is_internal: true,
5811 partitions: vec![PartitionMetadata {
5812 error_code: 0,
5813 partition_index: 0,
5814 leader_id: MetadataResponse::NO_LEADER_ID,
5815 leader_epoch: RecordBatch::NO_PARTITION_LEADER_EPOCH,
5816 replica_nodes: Vec::new(),
5817 isr_nodes: Vec::new(),
5818 offline_replicas: Vec::new(),
5819 }],
5820 topic_authorized_operations: 0,
5821 };
5822 assert_eq!(
5823 unnamed.to_string(),
5824 "TopicMetadata{error=UNKNOWN_TOPIC_OR_PARTITION, topic='null', topicId='AAAAAAAAAAAAAAAAAAAAAA', isInternal=true, partitionMetadata=[PartitionMetadata(error=NONE, partition=null-0, leader=Optional.empty, leaderEpoch=Optional.empty, replicas=, isr=, offlineReplicas=)], authorizedOperations=0}"
5825 );
5826 }
5827
5828 #[test]
5829 fn metadata_response_topic_metadata_matches_java() {
5830 assert_eq!(
5838 TopicMetadata::error(0, Some("t"), [0; 16]),
5839 TopicMetadata::new(0, "t", false, Vec::new())
5840 );
5841 let with = TopicMetadata::new(
5842 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
5843 "missing",
5844 true,
5845 Vec::new(),
5846 );
5847 assert_eq!(with.error_code, crate::error::UNKNOWN_TOPIC_OR_PARTITION);
5848 assert_eq!(with.name.as_deref(), Some("missing"));
5849 assert_eq!(with.topic_id, [0u8; 16]);
5850 assert!(with.is_internal);
5851 assert!(with.partitions.is_empty());
5852 assert_eq!(
5853 with.topic_authorized_operations,
5854 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
5855 );
5856 leftover_metadata_topic_metadata(1, std::slice::from_ref(&with));
5857 leftover_metadata_topic_metadata(1, &[]);
5858 leftover_metadata_topic_metadata(8, std::slice::from_ref(&with));
5859 leftover_metadata_topic_metadata(8, &[]);
5860 leftover_metadata_topic_metadata(10, std::slice::from_ref(&with));
5861 leftover_metadata_topic_metadata(10, &[]);
5862 leftover_metadata_topic_metadata(13, std::slice::from_ref(&with));
5863 leftover_metadata_topic_metadata(13, &[]);
5864 }
5865
5866 fn leftover_metadata_topic_metadata(version: i16, topics: &[TopicMetadata]) {
5867 let resp = MetadataResponse {
5868 throttle_time_ms: 0,
5869 brokers: Vec::new(),
5870 cluster_id: None,
5871 controller_id: 1,
5872 topics: topics.to_vec(),
5873 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
5874 error_code: 0,
5875 };
5876 let mut buf = BytesMut::new();
5877 encode_metadata_response(&mut buf, version, &resp).unwrap();
5878 let mut cur = buf.as_ref();
5879 let decoded = decode_metadata_response(&mut cur, version).unwrap();
5880 assert_eq!(decoded.throttle_time_ms, 0);
5881 assert!(decoded.brokers.is_empty());
5882 assert!(decoded.cluster_id.is_none());
5883 assert_eq!(decoded.controller_id, 1);
5884 assert_eq!(decoded.topics, topics);
5885 assert_eq!(
5886 decoded.cluster_authorized_operations,
5887 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
5888 );
5889 assert_eq!(decoded.error_code, 0);
5890 leftover_empty(
5891 &cur,
5892 match (version, topics.is_empty()) {
5893 (1, false) => "Metadata v1 TopicMetadata leftover-empty",
5894 (1, true) => "Metadata v1 TopicMetadata empty leftover-empty",
5895 (8, false) => "Metadata v8 TopicMetadata leftover-empty",
5896 (8, true) => "Metadata v8 TopicMetadata empty leftover-empty",
5897 (10, false) => "Metadata v10 TopicMetadata leftover-empty",
5898 (10, true) => "Metadata v10 TopicMetadata empty leftover-empty",
5899 (13, false) => "Metadata v13 TopicMetadata leftover-empty",
5900 (13, true) => "Metadata v13 TopicMetadata empty leftover-empty",
5901 _ => "Metadata TopicMetadata leftover-empty",
5902 },
5903 )
5904 .unwrap();
5905 }
5906
5907 #[test]
5908 fn metadata_response_topic_metadata_with_topic_id_matches_java() {
5909 let id = [7u8; 16];
5916 assert_eq!(
5917 TopicMetadata::new(0, "t", false, Vec::new()),
5918 TopicMetadata::with_topic_id(
5919 0,
5920 "t",
5921 [0; 16],
5922 false,
5923 Vec::new(),
5924 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
5925 )
5926 );
5927 let with = TopicMetadata::with_topic_id(
5928 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
5929 "missing",
5930 id,
5931 true,
5932 Vec::new(),
5933 0,
5934 );
5935 assert_eq!(with.error_code, crate::error::UNKNOWN_TOPIC_OR_PARTITION);
5936 assert_eq!(with.name.as_deref(), Some("missing"));
5937 assert_eq!(with.topic_id, id);
5938 assert!(with.is_internal);
5939 assert!(with.partitions.is_empty());
5940 assert_eq!(with.topic_authorized_operations, 0);
5941 leftover_metadata_topic_metadata_with_topic_id(1, std::slice::from_ref(&with));
5942 leftover_metadata_topic_metadata_with_topic_id(1, &[]);
5943 leftover_metadata_topic_metadata_with_topic_id(8, std::slice::from_ref(&with));
5944 leftover_metadata_topic_metadata_with_topic_id(8, &[]);
5945 leftover_metadata_topic_metadata_with_topic_id(10, std::slice::from_ref(&with));
5946 leftover_metadata_topic_metadata_with_topic_id(10, &[]);
5947 leftover_metadata_topic_metadata_with_topic_id(13, std::slice::from_ref(&with));
5948 leftover_metadata_topic_metadata_with_topic_id(13, &[]);
5949 }
5950
5951 fn leftover_metadata_topic_metadata_with_topic_id(version: i16, topics: &[TopicMetadata]) {
5952 let resp = MetadataResponse {
5953 throttle_time_ms: 0,
5954 brokers: Vec::new(),
5955 cluster_id: None,
5956 controller_id: 1,
5957 topics: topics.to_vec(),
5958 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
5959 error_code: 0,
5960 };
5961 let mut buf = BytesMut::new();
5962 encode_metadata_response(&mut buf, version, &resp).unwrap();
5963 let mut cur = buf.as_ref();
5964 let decoded = decode_metadata_response(&mut cur, version).unwrap();
5965 assert_eq!(decoded.throttle_time_ms, 0);
5966 assert!(decoded.brokers.is_empty());
5967 assert!(decoded.cluster_id.is_none());
5968 assert_eq!(decoded.controller_id, 1);
5969 if version >= 10 {
5970 assert_eq!(decoded.topics, topics);
5971 } else {
5972 assert_eq!(decoded.topics.len(), topics.len());
5973 for (got, want) in decoded.topics.iter().zip(topics) {
5974 assert_eq!(got.error_code, want.error_code);
5975 assert_eq!(got.name, want.name);
5976 assert_eq!(
5977 got.topic_id, [0u8; 16],
5978 "Metadata v{version} omits TopicId; decode fills zeros"
5979 );
5980 assert_eq!(got.is_internal, want.is_internal);
5981 assert_eq!(got.partitions, want.partitions);
5982 if version < 8 {
5983 assert_eq!(
5984 got.topic_authorized_operations,
5985 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
5986 "Metadata v{version} omits topicAuthorizedOperations; decode fills AUTHORIZED_OPERATIONS_OMITTED"
5987 );
5988 } else {
5989 assert_eq!(
5990 got.topic_authorized_operations,
5991 want.topic_authorized_operations
5992 );
5993 }
5994 }
5995 }
5996 assert_eq!(
5997 decoded.cluster_authorized_operations,
5998 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
5999 );
6000 assert_eq!(decoded.error_code, 0);
6001 leftover_empty(
6002 &cur,
6003 match (version, topics.is_empty()) {
6004 (1, false) => "Metadata v1 TopicMetadata.topicId leftover-empty",
6005 (1, true) => "Metadata v1 TopicMetadata.topicId empty leftover-empty",
6006 (8, false) => "Metadata v8 TopicMetadata.topicId leftover-empty",
6007 (8, true) => "Metadata v8 TopicMetadata.topicId empty leftover-empty",
6008 (10, false) => "Metadata v10 TopicMetadata.topicId leftover-empty",
6009 (10, true) => "Metadata v10 TopicMetadata.topicId empty leftover-empty",
6010 (13, false) => "Metadata v13 TopicMetadata.topicId leftover-empty",
6011 (13, true) => "Metadata v13 TopicMetadata.topicId empty leftover-empty",
6012 _ => "Metadata TopicMetadata.topicId leftover-empty",
6013 },
6014 )
6015 .unwrap();
6016 }
6017
6018 #[test]
6019 fn metadata_response_partition_metadata_matches_java() {
6020 let none = PartitionMetadata::new(0, 0, None, None, Vec::new(), Vec::new(), Vec::new());
6027 assert_eq!(none.error_code, 0);
6028 assert_eq!(none.partition_index, 0);
6029 assert_eq!(none.leader_id, MetadataResponse::NO_LEADER_ID);
6030 assert_eq!(none.leader_epoch, RecordBatch::NO_PARTITION_LEADER_EPOCH);
6031 assert!(none.replica_nodes.is_empty());
6032 assert!(none.isr_nodes.is_empty());
6033 assert!(none.offline_replicas.is_empty());
6034 let with = PartitionMetadata::new(0, 1, Some(2), Some(8), vec![2, 3], vec![2], vec![3]);
6035 assert_eq!(with.partition_index, 1);
6036 assert_eq!(with.leader_id, 2);
6037 assert_eq!(with.leader_epoch, 8);
6038 assert_eq!(with.replica_nodes, vec![2, 3]);
6039 assert_eq!(with.isr_nodes, vec![2]);
6040 assert_eq!(with.offline_replicas, vec![3]);
6041 assert_eq!(
6042 with.without_leader_epoch(),
6043 PartitionMetadata::new(0, 1, Some(2), None, vec![2, 3], vec![2], vec![3],)
6044 );
6045 leftover_metadata_partition_metadata(1, std::slice::from_ref(&with));
6046 leftover_metadata_partition_metadata(1, &[]);
6047 leftover_metadata_partition_metadata(5, std::slice::from_ref(&with));
6048 leftover_metadata_partition_metadata(5, &[]);
6049 leftover_metadata_partition_metadata(7, std::slice::from_ref(&with));
6050 leftover_metadata_partition_metadata(7, &[]);
6051 leftover_metadata_partition_metadata(13, std::slice::from_ref(&with));
6052 leftover_metadata_partition_metadata(13, &[]);
6053 }
6054
6055 fn leftover_metadata_partition_metadata(version: i16, parts: &[PartitionMetadata]) {
6056 let topics = if parts.is_empty() {
6057 Vec::new()
6058 } else {
6059 vec![TopicMetadata::new(0, "t", false, parts.to_vec())]
6060 };
6061 let resp = MetadataResponse {
6062 throttle_time_ms: 0,
6063 brokers: Vec::new(),
6064 cluster_id: None,
6065 controller_id: 1,
6066 topics,
6067 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6068 error_code: 0,
6069 };
6070 let mut buf = BytesMut::new();
6071 encode_metadata_response(&mut buf, version, &resp).unwrap();
6072 let mut cur = buf.as_ref();
6073 let decoded = decode_metadata_response(&mut cur, version).unwrap();
6074 assert_eq!(decoded.throttle_time_ms, 0);
6075 assert!(decoded.brokers.is_empty());
6076 assert!(decoded.cluster_id.is_none());
6077 assert_eq!(decoded.controller_id, 1);
6078 if version >= 7 {
6079 assert_eq!(decoded.topics, resp.topics);
6080 } else {
6081 assert_eq!(decoded.topics.len(), resp.topics.len());
6082 for (got_topic, want_topic) in decoded.topics.iter().zip(&resp.topics) {
6083 assert_eq!(got_topic.error_code, want_topic.error_code);
6084 assert_eq!(got_topic.name, want_topic.name);
6085 assert_eq!(got_topic.topic_id, [0u8; 16]);
6086 assert_eq!(got_topic.is_internal, want_topic.is_internal);
6087 assert_eq!(got_topic.partitions.len(), want_topic.partitions.len());
6088 for (got, want) in got_topic.partitions.iter().zip(&want_topic.partitions) {
6089 assert_eq!(got.error_code, want.error_code);
6090 assert_eq!(got.partition_index, want.partition_index);
6091 assert_eq!(got.leader_id, want.leader_id);
6092 assert_eq!(
6093 got.leader_epoch,
6094 RecordBatch::NO_PARTITION_LEADER_EPOCH,
6095 "Metadata v{version} omits LeaderEpoch; decode fills NO_PARTITION_LEADER_EPOCH"
6096 );
6097 assert_eq!(got.replica_nodes, want.replica_nodes);
6098 assert_eq!(got.isr_nodes, want.isr_nodes);
6099 if version < 5 {
6100 assert!(
6101 got.offline_replicas.is_empty(),
6102 "Metadata v{version} omits OfflineReplicas; decode fills empty"
6103 );
6104 } else {
6105 assert_eq!(got.offline_replicas, want.offline_replicas);
6106 }
6107 }
6108 assert_eq!(
6109 got_topic.topic_authorized_operations,
6110 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
6111 );
6112 }
6113 }
6114 assert_eq!(
6115 decoded.cluster_authorized_operations,
6116 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
6117 );
6118 assert_eq!(decoded.error_code, 0);
6119 leftover_empty(
6120 &cur,
6121 match (version, parts.is_empty()) {
6122 (1, false) => "Metadata v1 PartitionMetadata leftover-empty",
6123 (1, true) => "Metadata v1 PartitionMetadata empty leftover-empty",
6124 (5, false) => "Metadata v5 PartitionMetadata leftover-empty",
6125 (5, true) => "Metadata v5 PartitionMetadata empty leftover-empty",
6126 (7, false) => "Metadata v7 PartitionMetadata leftover-empty",
6127 (7, true) => "Metadata v7 PartitionMetadata empty leftover-empty",
6128 (13, false) => "Metadata v13 PartitionMetadata leftover-empty",
6129 (13, true) => "Metadata v13 PartitionMetadata empty leftover-empty",
6130 _ => "Metadata PartitionMetadata leftover-empty",
6131 },
6132 )
6133 .unwrap();
6134 }
6135
6136 #[test]
6137 fn metadata_response_errors_matches_java() {
6138 fn topic(error_code: i16, name: Option<&str>, topic_id: [u8; 16]) -> TopicMetadata {
6139 TopicMetadata {
6140 error_code,
6141 name: name.map(str::to_owned),
6142 topic_id,
6143 is_internal: false,
6144 partitions: Vec::new(),
6145 topic_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6146 }
6147 }
6148 fn resp(topics: Vec<TopicMetadata>) -> MetadataResponse {
6149 MetadataResponse {
6150 throttle_time_ms: 0,
6151 brokers: Vec::new(),
6152 cluster_id: None,
6153 controller_id: MetadataResponse::NO_CONTROLLER_ID,
6154 topics,
6155 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6156 error_code: 0,
6157 }
6158 }
6159
6160 let named = resp(vec![
6161 topic(0, Some("ok"), [0; 16]),
6162 topic(
6163 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
6164 Some("missing"),
6165 [0; 16],
6166 ),
6167 topic(
6168 crate::error::TOPIC_AUTHORIZATION_FAILED,
6169 Some("denied"),
6170 [0; 16],
6171 ),
6172 topic(crate::error::INVALID_TOPIC_EXCEPTION, Some("bad"), [0; 16]),
6173 ]);
6174 assert_eq!(
6175 named.errors().unwrap(),
6176 HashMap::from([
6177 (
6178 "missing".to_owned(),
6179 crate::error::UNKNOWN_TOPIC_OR_PARTITION
6180 ),
6181 (
6182 "denied".to_owned(),
6183 crate::error::TOPIC_AUTHORIZATION_FAILED
6184 ),
6185 ("bad".to_owned(), crate::error::INVALID_TOPIC_EXCEPTION),
6186 ])
6187 );
6188 assert_eq!(
6189 named.topics_by_error(crate::error::UNKNOWN_TOPIC_OR_PARTITION),
6190 HashSet::from(["missing".to_owned()])
6191 );
6192 assert_eq!(named.topics_by_error(0), HashSet::from(["ok".to_owned()]));
6193 assert!(named
6194 .topics_by_error(crate::error::UNKNOWN_TOPIC_ID)
6195 .is_empty());
6196 let err = named.errors_by_topic_id().unwrap_err();
6197 assert!(
6198 err.to_string()
6199 .contains("Use errors() when managing topic using topic name"),
6200 "got {err}"
6201 );
6202
6203 let by_id = resp(vec![
6204 topic(0, None, [1; 16]),
6205 topic(crate::error::UNKNOWN_TOPIC_ID, None, [2; 16]),
6206 topic(
6207 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
6208 Some("named"),
6209 [3; 16],
6210 ),
6211 ]);
6212 assert_eq!(
6213 by_id.errors_by_topic_id().unwrap(),
6214 HashMap::from([
6215 ([2; 16], crate::error::UNKNOWN_TOPIC_ID),
6216 ([3; 16], crate::error::UNKNOWN_TOPIC_OR_PARTITION),
6217 ])
6218 );
6219 assert_eq!(
6220 by_id.topics_by_error(crate::error::UNKNOWN_TOPIC_ID),
6221 HashSet::new()
6222 );
6223 assert_eq!(
6224 by_id.topics_by_error(crate::error::UNKNOWN_TOPIC_OR_PARTITION),
6225 HashSet::from(["named".to_owned()])
6226 );
6227 let err = by_id.errors().unwrap_err();
6228 assert!(
6229 err.to_string()
6230 .contains("Use errorsByTopicId() when managing topic using topic id"),
6231 "got {err}"
6232 );
6233
6234 let mixed_null_name = resp(vec![
6235 topic(crate::error::UNKNOWN_TOPIC_OR_PARTITION, Some("t"), [1; 16]),
6236 topic(0, None, [2; 16]),
6237 ]);
6238 let err = mixed_null_name.errors().unwrap_err();
6239 assert!(
6240 err.to_string()
6241 .contains("Use errorsByTopicId() when managing topic using topic id"),
6242 "got {err}"
6243 );
6244
6245 let mixed_zero_id = resp(vec![
6246 topic(crate::error::UNKNOWN_TOPIC_ID, Some("t"), [1; 16]),
6247 topic(0, Some("ok"), [0; 16]),
6248 ]);
6249 let err = mixed_zero_id.errors_by_topic_id().unwrap_err();
6250 assert!(
6251 err.to_string()
6252 .contains("Use errors() when managing topic using topic name"),
6253 "got {err}"
6254 );
6255
6256 let empty = resp(Vec::new());
6257 assert!(empty.errors().unwrap().is_empty());
6258 assert!(empty.errors_by_topic_id().unwrap().is_empty());
6259 assert!(empty.topics_by_error(0).is_empty());
6260 assert!(empty.error_counts().is_empty());
6261 assert!(empty.topic_authorized_operations("t").is_none());
6262 }
6263
6264 #[test]
6265 fn metadata_response_error_counts_matches_java() {
6266 fn topic(
6267 error_code: i16,
6268 name: Option<&str>,
6269 partitions: Vec<i16>,
6270 authorized: i32,
6271 ) -> TopicMetadata {
6272 TopicMetadata {
6273 error_code,
6274 name: name.map(str::to_owned),
6275 topic_id: [0; 16],
6276 is_internal: false,
6277 partitions: partitions
6278 .into_iter()
6279 .enumerate()
6280 .map(|(i, error_code)| PartitionMetadata {
6281 error_code,
6282 partition_index: i32::try_from(i).unwrap_or(i32::MAX),
6283 leader_id: MetadataResponse::NO_LEADER_ID,
6284 leader_epoch: RecordBatch::NO_PARTITION_LEADER_EPOCH,
6285 replica_nodes: Vec::new(),
6286 isr_nodes: Vec::new(),
6287 offline_replicas: Vec::new(),
6288 })
6289 .collect(),
6290 topic_authorized_operations: authorized,
6291 }
6292 }
6293 let resp = MetadataResponse {
6294 throttle_time_ms: 0,
6295 brokers: Vec::new(),
6296 cluster_id: None,
6297 controller_id: MetadataResponse::NO_CONTROLLER_ID,
6298 topics: vec![
6299 topic(
6300 0,
6301 Some("ok"),
6302 vec![0, crate::error::NOT_LEADER_OR_FOLLOWER],
6303 1,
6304 ),
6305 topic(
6306 crate::error::UNKNOWN_TOPIC_OR_PARTITION,
6307 Some("missing"),
6308 Vec::new(),
6309 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6310 ),
6311 topic(0, None, vec![0], 7),
6312 ],
6313 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6314 error_code: crate::error::TOPIC_AUTHORIZATION_FAILED,
6315 };
6316 assert_eq!(
6317 resp.error_counts(),
6318 HashMap::from([
6319 (0, 4),
6320 (crate::error::NOT_LEADER_OR_FOLLOWER, 1),
6321 (crate::error::UNKNOWN_TOPIC_OR_PARTITION, 1),
6322 ])
6323 );
6324 assert_eq!(resp.topic_authorized_operations("ok"), Some(1));
6325 assert_eq!(
6326 resp.topic_authorized_operations("missing"),
6327 Some(MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED)
6328 );
6329 assert!(resp.topic_authorized_operations("nope").is_none());
6330 assert!(resp.topic_authorized_operations("").is_none());
6331 }
6332
6333 #[test]
6334 fn metadata_response_brokers_by_id_matches_java() {
6335 let a = Broker::new(1, "127.0.0.1", 9092, Some("r".into()));
6336 let b = Broker::new(3, "10.0.0.3", 9093, None);
6337 let resp = MetadataResponse {
6338 throttle_time_ms: 0,
6339 brokers: vec![a.clone(), b.clone()],
6340 cluster_id: None,
6341 controller_id: 1,
6342 topics: Vec::new(),
6343 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6344 error_code: 0,
6345 };
6346 assert_eq!(resp.brokers_by_id(), HashMap::from([(1, a), (3, b)]));
6347 let empty = MetadataResponse {
6348 throttle_time_ms: 0,
6349 brokers: Vec::new(),
6350 cluster_id: None,
6351 controller_id: MetadataResponse::NO_CONTROLLER_ID,
6352 topics: Vec::new(),
6353 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6354 error_code: 0,
6355 };
6356 assert!(empty.brokers_by_id().is_empty());
6357 }
6358
6359 #[test]
6360 fn metadata_response_controller_matches_java() {
6361 let a = Broker::new(1, "127.0.0.1", 9092, Some("r".into()));
6366 let b = Broker::new(3, "10.0.0.3", 9093, None);
6367 let resp = MetadataResponse {
6368 throttle_time_ms: 0,
6369 brokers: vec![a.clone(), b.clone()],
6370 cluster_id: None,
6371 controller_id: 1,
6372 topics: Vec::new(),
6373 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6374 error_code: 0,
6375 };
6376 assert_eq!(resp.controller(), Some(&a));
6377 let other = MetadataResponse {
6378 controller_id: 3,
6379 ..resp.clone()
6380 };
6381 assert_eq!(other.controller(), Some(&b));
6382 let missing = MetadataResponse {
6383 controller_id: 2,
6384 ..resp.clone()
6385 };
6386 assert!(missing.controller().is_none());
6387 let empty = MetadataResponse {
6388 throttle_time_ms: 0,
6389 brokers: Vec::new(),
6390 cluster_id: None,
6391 controller_id: MetadataResponse::NO_CONTROLLER_ID,
6392 topics: Vec::new(),
6393 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6394 error_code: 0,
6395 };
6396 assert!(empty.controller().is_none());
6397 let dup_last = Broker::new(1, "10.0.0.9", 9094, None);
6398 let dups = MetadataResponse {
6399 brokers: vec![a.clone(), dup_last.clone()],
6400 controller_id: 1,
6401 ..empty.clone()
6402 };
6403 assert_eq!(dups.controller(), Some(&dup_last));
6404 leftover_metadata_controller(1, &resp);
6405 leftover_metadata_controller(1, &empty);
6406 leftover_metadata_controller(8, &resp);
6407 leftover_metadata_controller(8, &empty);
6408 leftover_metadata_controller(13, &resp);
6409 leftover_metadata_controller(13, &empty);
6410 }
6411
6412 fn leftover_metadata_controller(version: i16, resp: &MetadataResponse) {
6413 let mut buf = BytesMut::new();
6414 encode_metadata_response(&mut buf, version, resp).unwrap();
6415 let mut cur = buf.as_ref();
6416 let decoded = decode_metadata_response(&mut cur, version).unwrap();
6417 assert_eq!(decoded, *resp);
6418 assert_eq!(decoded.controller(), resp.controller());
6419 leftover_empty(
6420 &cur,
6421 match (version, resp.brokers.is_empty()) {
6422 (1, false) => "Metadata v1 controller leftover-empty",
6423 (1, true) => "Metadata v1 controller empty leftover-empty",
6424 (8, false) => "Metadata v8 controller leftover-empty",
6425 (8, true) => "Metadata v8 controller empty leftover-empty",
6426 (13, false) => "Metadata v13 controller leftover-empty",
6427 (13, true) => "Metadata v13 controller empty leftover-empty",
6428 _ => "Metadata controller leftover-empty",
6429 },
6430 )
6431 .unwrap();
6432 }
6433
6434 #[test]
6435 fn metadata_v7_decodes_leader_epoch_and_omitted_authorized_ops() {
6436 let resp = MetadataResponse {
6437 throttle_time_ms: 0,
6438 brokers: vec![Broker {
6439 node_id: 1,
6440 host: "127.0.0.1".into(),
6441 port: 9092,
6442 rack: None,
6443 }],
6444 cluster_id: Some("cid".into()),
6445 controller_id: 1,
6446 topics: vec![TopicMetadata {
6447 error_code: 0,
6448 name: Some("orders".into()),
6449 topic_id: [0u8; 16],
6450 is_internal: false,
6451 partitions: vec![PartitionMetadata {
6452 error_code: 0,
6453 partition_index: 0,
6454 leader_id: 1,
6455 leader_epoch: 3,
6456 replica_nodes: vec![1],
6457 isr_nodes: vec![1],
6458 offline_replicas: Vec::new(),
6459 }],
6460 topic_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6461 }],
6462 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
6463 error_code: 0,
6464 };
6465 let mut buf = BytesMut::new();
6466 encode_metadata_response(&mut buf, 7, &resp).unwrap();
6467 let mut cur = &buf[..];
6468 let decoded = decode_metadata_response(&mut cur, 7).unwrap();
6469 leftover_empty(&cur, "Metadata v7").unwrap();
6470 assert_eq!(decoded, resp);
6471 assert_eq!(decoded.topics[0].partitions[0].leader_epoch, 3);
6472 assert_eq!(
6473 decoded.topics[0].topic_authorized_operations,
6474 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
6475 );
6476 assert!(
6477 !MetadataResponse::has_reliable_leader_epochs(7),
6478 "v7 leader epochs are on the wire but must not be retained by the client"
6479 );
6480 }
6481
6482 #[test]
6483 fn metadata_request_allow_auto_topic_creation_matches_java() {
6484 leftover_metadata_allow_auto_topic_creation(1, true);
6493 leftover_metadata_allow_auto_topic_creation(3, true);
6494 leftover_metadata_allow_auto_topic_creation(4, true);
6495 leftover_metadata_allow_auto_topic_creation(4, false);
6496 leftover_metadata_allow_auto_topic_creation(8, true);
6497 leftover_metadata_allow_auto_topic_creation(13, true);
6498 leftover_metadata_allow_auto_topic_creation(13, false);
6499
6500 let mut v3 = BytesMut::new();
6501 encode_metadata_request(&mut v3, 3, None, true).unwrap();
6502 let mut v4_true = BytesMut::new();
6503 encode_metadata_request(&mut v4_true, 4, None, true).unwrap();
6504 assert_ne!(
6505 &v3[..],
6506 &v4_true[..],
6507 "v4 writes AllowAutoTopicCreation after Topics"
6508 );
6509 let mut v4_false = BytesMut::new();
6510 encode_metadata_request(&mut v4_false, 4, None, false).unwrap();
6511 assert_ne!(
6512 &v4_true[..],
6513 &v4_false[..],
6514 "v4 AllowAutoTopicCreation is not always the JSON default true"
6515 );
6516 let mut v1 = BytesMut::new();
6517 encode_metadata_request(&mut v1, 1, None, true).unwrap();
6518 assert_eq!(
6519 &v1[..],
6520 &v3[..],
6521 "v1 and v3 both omit AllowAutoTopicCreation"
6522 );
6523 }
6524
6525 #[test]
6526 fn metadata_request_all_topics_matches_java() {
6527 let (topics, allow) = MetadataRequest::all_topics();
6534 assert!(topics.is_none());
6535 assert!(allow);
6536 assert!(MetadataRequest::is_all_topics(1, topics));
6537 assert!(MetadataRequest::is_all_topics(13, topics));
6538 leftover_metadata_all_topics(1);
6539 leftover_metadata_all_topics(3);
6540 leftover_metadata_all_topics(4);
6541 leftover_metadata_all_topics(9);
6542 leftover_metadata_all_topics(13);
6543
6544 let mut all = BytesMut::new();
6545 encode_metadata_request(&mut all, 4, None, allow).unwrap();
6546 let named = ["t".to_string()];
6547 let mut with_names = BytesMut::new();
6548 encode_metadata_request(&mut with_names, 4, Some(&named), allow).unwrap();
6549 assert_ne!(
6550 &all[..],
6551 &with_names[..],
6552 "allTopics null Topics is not an empty or named Topics list"
6553 );
6554 }
6555
6556 #[test]
6557 fn metadata_request_for_topic_ids_matches_java() {
6558 let (none, allow) = MetadataRequest::for_topic_ids(None);
6567 assert!(none.is_none());
6568 assert!(!allow);
6569 assert!(MetadataRequest::is_all_topics(1, none.as_deref()));
6570 assert_ne!(
6571 MetadataRequest::all_topics().1,
6572 allow,
6573 "allTopics AllowAutoTopicCreation is true"
6574 );
6575 let (empty, empty_allow) = MetadataRequest::for_topic_ids(Some(&[]));
6576 assert_eq!(empty.as_deref(), Some(&[][..]));
6577 assert!(!empty_allow);
6578 assert!(!MetadataRequest::is_all_topics(1, empty.as_deref()));
6579 let id = [7u8; 16];
6580 let (named, named_allow) = MetadataRequest::for_topic_ids(Some(std::slice::from_ref(&id)));
6581 assert_eq!(named, Some(MetadataRequestTopic::convert_from_ids([id])));
6582 assert!(!named_allow);
6583 leftover_metadata_for_topic_ids(4, None);
6584 leftover_metadata_for_topic_ids(13, None);
6585 leftover_metadata_for_topic_ids(4, Some(&[]));
6586 leftover_metadata_for_topic_ids(13, Some(&[]));
6587 leftover_metadata_for_topic_ids(12, Some(std::slice::from_ref(&id)));
6588 leftover_metadata_for_topic_ids(13, Some(std::slice::from_ref(&id)));
6589 }
6590
6591 fn leftover_metadata_for_topic_ids(version: i16, ids: Option<&[[u8; 16]]>) {
6592 let (topics, allow) = MetadataRequest::for_topic_ids(ids);
6593 assert!(!allow);
6594 let mut buf = BytesMut::new();
6595 encode_metadata_request_topics(&mut buf, version, topics.as_deref(), allow, false).unwrap();
6596 let mut cur = buf.as_ref();
6597 let (got, allow_got, include_topic, include_cluster) =
6598 decode_metadata_request_topics(&mut cur, version).unwrap();
6599 assert_eq!(got, topics);
6600 assert!(!allow_got);
6601 assert!(!include_topic);
6602 assert!(!include_cluster);
6603 let empty = if ids.is_some_and(<[[u8; 16]]>::is_empty) {
6604 "empty "
6605 } else {
6606 ""
6607 };
6608 assert!(
6609 cur.is_empty(),
6610 "Metadata v{version} Builder.topicIds {empty}leftover-empty; leftover {} bytes",
6611 cur.len()
6612 );
6613 }
6614
6615 #[test]
6616 fn metadata_request_for_topic_names_matches_java() {
6617 let (none_true, allow_true) = MetadataRequest::for_topic_names(None, true);
6627 assert!(none_true.is_none());
6628 assert!(allow_true);
6629 assert_eq!(allow_true, MetadataRequest::all_topics().1);
6630 assert!(MetadataRequest::is_all_topics(1, none_true.as_deref()));
6631 let (none_false, allow_false) = MetadataRequest::for_topic_names(None, false);
6632 assert!(none_false.is_none());
6633 assert!(!allow_false);
6634 assert_eq!(allow_false, MetadataRequest::for_topic_ids(None).1);
6635 let names = ["t".to_string()];
6636 let (named, named_allow) = MetadataRequest::for_topic_names(Some(&names), true);
6637 assert_eq!(named, Some(MetadataRequestTopic::convert_from_names(["t"])));
6638 assert!(named_allow);
6639 let (empty, empty_allow) = MetadataRequest::for_topic_names(Some(&[]), false);
6640 assert_eq!(empty.as_deref(), Some(&[][..]));
6641 assert!(!empty_allow);
6642 leftover_metadata_for_topic_names(1, None, true);
6643 leftover_metadata_for_topic_names(13, None, true);
6644 leftover_metadata_for_topic_names(4, None, false);
6645 leftover_metadata_for_topic_names(13, None, false);
6646 leftover_metadata_for_topic_names(1, Some(&names), true);
6647 leftover_metadata_for_topic_names(13, Some(&names), true);
6648 leftover_metadata_for_topic_names(4, Some(&names), false);
6649 leftover_metadata_for_topic_names(13, Some(&names), false);
6650 leftover_metadata_for_topic_names(1, Some(&[]), true);
6651 leftover_metadata_for_topic_names(13, Some(&[]), true);
6652 leftover_metadata_for_topic_names(4, Some(&[]), false);
6653 leftover_metadata_for_topic_names(13, Some(&[]), false);
6654 }
6655
6656 #[test]
6657 fn metadata_request_for_topic_names_version_matches_java() {
6658 let names = ["t".to_string()];
6668 let (none, allow, oldest, latest) =
6669 MetadataRequest::for_topic_names_version(None, true, 13);
6670 assert!(none.is_none());
6671 assert!(allow);
6672 assert_eq!(oldest, 13);
6673 assert_eq!(latest, 13);
6674 assert_eq!(allow, MetadataRequest::all_topics().1);
6675 let (named, named_allow, named_oldest, named_latest) =
6676 MetadataRequest::for_topic_names_version(Some(&names), false, 4);
6677 assert_eq!(named, Some(MetadataRequestTopic::convert_from_names(["t"])));
6678 assert!(!named_allow);
6679 assert_eq!(named_oldest, 4);
6680 assert_eq!(named_latest, 4);
6681 assert_eq!(
6682 MetadataRequest::for_topic_names_version(None, false, 1),
6683 (None, false, 1, 1)
6684 );
6685 let (from_names, from_allow) = MetadataRequest::for_topic_names(Some(&names), false);
6686 assert_eq!(named, from_names);
6687 assert_eq!(named_allow, from_allow);
6688 leftover_metadata_for_topic_names_version(1, None, true, 13);
6689 leftover_metadata_for_topic_names_version(13, None, true, 13);
6690 leftover_metadata_for_topic_names_version(4, None, false, 4);
6691 leftover_metadata_for_topic_names_version(13, None, false, 4);
6692 leftover_metadata_for_topic_names_version(1, Some(&names), true, 13);
6693 leftover_metadata_for_topic_names_version(13, Some(&names), true, 13);
6694 leftover_metadata_for_topic_names_version(4, Some(&names), false, 4);
6695 leftover_metadata_for_topic_names_version(13, Some(&names), false, 4);
6696 leftover_metadata_for_topic_names_version(1, Some(&[]), true, 1);
6697 leftover_metadata_for_topic_names_version(13, Some(&[]), true, 1);
6698 leftover_metadata_for_topic_names_version(4, Some(&[]), false, 12);
6699 leftover_metadata_for_topic_names_version(13, Some(&[]), false, 12);
6700 }
6701
6702 #[test]
6703 fn metadata_request_for_topic_names_range_matches_java() {
6704 let names = ["t".to_string()];
6715 let (none, allow, oldest, latest) =
6716 MetadataRequest::for_topic_names_range(None, true, 1, 13);
6717 assert!(none.is_none());
6718 assert!(allow);
6719 assert_eq!(oldest, 1);
6720 assert_eq!(latest, 13);
6721 assert_eq!(allow, MetadataRequest::all_topics().1);
6722 let (named, named_allow, named_oldest, named_latest) =
6723 MetadataRequest::for_topic_names_range(Some(&names), false, 4, 13);
6724 assert_eq!(named, Some(MetadataRequestTopic::convert_from_names(["t"])));
6725 assert!(!named_allow);
6726 assert_eq!(named_oldest, 4);
6727 assert_eq!(named_latest, 13);
6728 assert_eq!(
6729 MetadataRequest::for_topic_names_version(Some(&names), false, 4),
6730 MetadataRequest::for_topic_names_range(Some(&names), false, 4, 4)
6731 );
6732 assert_eq!(
6733 MetadataRequest::for_topic_names_range(None, false, 13, 1),
6734 (None, false, 13, 1)
6735 );
6736 leftover_metadata_for_topic_names_range(1, None, true, 1, 13);
6737 leftover_metadata_for_topic_names_range(13, None, true, 1, 13);
6738 leftover_metadata_for_topic_names_range(4, None, false, 4, 13);
6739 leftover_metadata_for_topic_names_range(13, None, false, 4, 13);
6740 leftover_metadata_for_topic_names_range(1, Some(&names), true, 1, 13);
6741 leftover_metadata_for_topic_names_range(13, Some(&names), true, 1, 13);
6742 leftover_metadata_for_topic_names_range(4, Some(&names), false, 4, 13);
6743 leftover_metadata_for_topic_names_range(13, Some(&names), false, 4, 13);
6744 leftover_metadata_for_topic_names_range(1, Some(&[]), true, 1, 1);
6745 leftover_metadata_for_topic_names_range(13, Some(&[]), true, 1, 1);
6746 leftover_metadata_for_topic_names_range(4, Some(&[]), false, 4, 12);
6747 leftover_metadata_for_topic_names_range(13, Some(&[]), false, 4, 12);
6748 }
6749
6750 #[test]
6751 fn metadata_request_build_matches_java() {
6752 let named = [MetadataRequestTopic::by_name("t")];
6763 MetadataRequest::build(1, Some(&named), true).unwrap();
6764 MetadataRequest::build(1, None, true).unwrap();
6765 MetadataRequest::build(4, Some(&named), false).unwrap();
6766 let v0 = MetadataRequest::build(0, Some(&named), true).unwrap_err();
6767 assert!(
6768 matches!(v0, Error::Unsupported(_)),
6769 "v0 is Java UnsupportedVersionException, got {v0}"
6770 );
6771 assert!(v0.to_string().contains("older than 1"), "got {v0}");
6772 let auto = MetadataRequest::build(3, Some(&named), false).unwrap_err();
6773 assert!(
6774 matches!(auto, Error::Unsupported(_)),
6775 "allowAuto false below v4 is Java UnsupportedVersionException, got {auto}"
6776 );
6777 assert!(
6778 auto.to_string().contains("allowAutoTopicCreation"),
6779 "got {auto}"
6780 );
6781 let id = [MetadataRequestTopic::by_id([1; 16])];
6782 let null_name = MetadataRequest::build(11, Some(&id), true).unwrap_err();
6783 assert!(
6784 matches!(null_name, Error::Unsupported(_)),
6785 "null Name below v12 is Java UnsupportedVersionException, got {null_name}"
6786 );
6787 assert!(
6788 null_name.to_string().contains("null topic names"),
6789 "got {null_name}"
6790 );
6791 let named_with_id = [MetadataRequestTopic {
6792 name: Some("t".into()),
6793 topic_id: [1; 16],
6794 }];
6795 let nonzero = MetadataRequest::build(11, Some(&named_with_id), true).unwrap_err();
6796 assert!(
6797 matches!(nonzero, Error::Unsupported(_)),
6798 "non-zero TopicId below v12 is Java UnsupportedVersionException, got {nonzero}"
6799 );
6800 assert!(
6801 nonzero.to_string().contains("non-zero topic IDs"),
6802 "got {nonzero}"
6803 );
6804 MetadataRequest::build(12, Some(&id), true).unwrap();
6805 MetadataRequest::build(12, Some(&named_with_id), false).unwrap();
6806 leftover_metadata_build(1, Some(&named), true);
6807 leftover_metadata_build(1, None, true);
6808 leftover_metadata_build(1, Some(&[]), true);
6809 leftover_metadata_build(4, Some(&named), false);
6810 leftover_metadata_build(12, Some(&id), true);
6811 leftover_metadata_build(13, Some(&named), true);
6812 }
6813
6814 fn leftover_metadata_build(
6815 version: i16,
6816 topics: Option<&[MetadataRequestTopic]>,
6817 allow_auto: bool,
6818 ) {
6819 MetadataRequest::build(version, topics, allow_auto).unwrap();
6820 let mut buf = BytesMut::new();
6821 encode_metadata_request_topics(&mut buf, version, topics, allow_auto, false).unwrap();
6822 let mut cur = buf.as_ref();
6823 let (got, allow_got, include_topic, include_cluster) =
6824 decode_metadata_request_topics(&mut cur, version).unwrap();
6825 assert_eq!(got.as_deref(), topics);
6826 assert_eq!(allow_got, allow_auto);
6827 assert!(!include_topic);
6828 assert!(!include_cluster);
6829 let empty = if topics.is_some_and(<[MetadataRequestTopic]>::is_empty) {
6830 "empty "
6831 } else {
6832 ""
6833 };
6834 assert!(
6835 cur.is_empty(),
6836 "Metadata v{version} Builder.build {empty}leftover-empty; leftover {} bytes",
6837 cur.len()
6838 );
6839 }
6840
6841 fn leftover_metadata_for_topic_names(version: i16, names: Option<&[String]>, allow_auto: bool) {
6842 let (topics, allow) = MetadataRequest::for_topic_names(names, allow_auto);
6843 assert_eq!(allow, allow_auto);
6844 let mut buf = BytesMut::new();
6845 encode_metadata_request_topics(&mut buf, version, topics.as_deref(), allow, false).unwrap();
6846 let mut cur = buf.as_ref();
6847 let (got, allow_got, include_topic, include_cluster) =
6848 decode_metadata_request_topics(&mut cur, version).unwrap();
6849 assert_eq!(got, topics);
6850 assert_eq!(allow_got, allow_auto);
6851 assert!(!include_topic);
6852 assert!(!include_cluster);
6853 let empty = if names.is_some_and(<[String]>::is_empty) {
6854 "empty "
6855 } else {
6856 ""
6857 };
6858 assert!(
6859 cur.is_empty(),
6860 "Metadata v{version} Builder.topicNames {empty}leftover-empty; leftover {} bytes",
6861 cur.len()
6862 );
6863 }
6864
6865 fn leftover_metadata_for_topic_names_version(
6866 version: i16,
6867 names: Option<&[String]>,
6868 allow_auto: bool,
6869 allowed_version: i16,
6870 ) {
6871 let (topics, allow, oldest, latest) =
6872 MetadataRequest::for_topic_names_version(names, allow_auto, allowed_version);
6873 assert_eq!(allow, allow_auto);
6874 assert_eq!(oldest, allowed_version);
6875 assert_eq!(latest, allowed_version);
6876 let mut buf = BytesMut::new();
6877 encode_metadata_request_topics(&mut buf, version, topics.as_deref(), allow, false).unwrap();
6878 let mut cur = buf.as_ref();
6879 let (got, allow_got, include_topic, include_cluster) =
6880 decode_metadata_request_topics(&mut cur, version).unwrap();
6881 assert_eq!(got, topics);
6882 assert_eq!(allow_got, allow_auto);
6883 assert!(!include_topic);
6884 assert!(!include_cluster);
6885 let empty = if names.is_some_and(<[String]>::is_empty) {
6886 "empty "
6887 } else {
6888 ""
6889 };
6890 assert!(
6891 cur.is_empty(),
6892 "Metadata v{version} Builder.topicNames.allowedVersion {empty}leftover-empty; leftover {} bytes",
6893 cur.len()
6894 );
6895 }
6896
6897 fn leftover_metadata_for_topic_names_range(
6898 version: i16,
6899 names: Option<&[String]>,
6900 allow_auto: bool,
6901 min_version: i16,
6902 max_version: i16,
6903 ) {
6904 let (topics, allow, oldest, latest) =
6905 MetadataRequest::for_topic_names_range(names, allow_auto, min_version, max_version);
6906 assert_eq!(allow, allow_auto);
6907 assert_eq!(oldest, min_version);
6908 assert_eq!(latest, max_version);
6909 let mut buf = BytesMut::new();
6910 encode_metadata_request_topics(&mut buf, version, topics.as_deref(), allow, false).unwrap();
6911 let mut cur = buf.as_ref();
6912 let (got, allow_got, include_topic, include_cluster) =
6913 decode_metadata_request_topics(&mut cur, version).unwrap();
6914 assert_eq!(got, topics);
6915 assert_eq!(allow_got, allow_auto);
6916 assert!(!include_topic);
6917 assert!(!include_cluster);
6918 let empty = if names.is_some_and(<[String]>::is_empty) {
6919 "empty "
6920 } else {
6921 ""
6922 };
6923 assert!(
6924 cur.is_empty(),
6925 "Metadata v{version} Builder.topicNames.minVersion.maxVersion {empty}leftover-empty; leftover {} bytes",
6926 cur.len()
6927 );
6928 }
6929
6930 fn leftover_metadata_all_topics(version: i16) {
6931 let (topics, allow) = MetadataRequest::all_topics();
6932 assert!(topics.is_none());
6933 let mut buf = BytesMut::new();
6934 encode_metadata_request(&mut buf, version, None, allow).unwrap();
6935 let mut cur = buf.as_ref();
6936 let (got, allow_got, include_topic, include_cluster) =
6937 decode_metadata_request_topics(&mut cur, version).unwrap();
6938 assert!(got.is_none());
6939 assert!(allow_got);
6940 assert!(!include_topic);
6941 assert!(!include_cluster);
6942 assert!(
6943 cur.is_empty(),
6944 "Metadata v{version} Builder.allTopics leftover-empty; leftover {} bytes",
6945 cur.len()
6946 );
6947 }
6948
6949 fn leftover_metadata_allow_auto_topic_creation(version: i16, allow_auto: bool) {
6950 let mut buf = BytesMut::new();
6951 encode_metadata_request(&mut buf, version, None, allow_auto).unwrap();
6952 let mut cur = buf.as_ref();
6953 let (topics, allow, include_topic, include_cluster) =
6954 decode_metadata_request_topics(&mut cur, version).unwrap();
6955 assert!(topics.is_none(), "all-topics is a null Topics array");
6956 let expected = if version >= 4 { allow_auto } else { true };
6957 assert_eq!(allow, expected);
6958 assert!(!include_topic);
6959 assert!(!include_cluster);
6960 let empty = if allow_auto { "empty " } else { "" };
6961 assert!(
6962 cur.is_empty(),
6963 "Metadata v{version} AllowAutoTopicCreation {empty}leftover-empty; leftover {} bytes",
6964 cur.len()
6965 );
6966 }
6967
6968 #[test]
6969 fn metadata_request_roundtrips_topics_and_allow_auto() {
6970 let topics = ["orders".to_string(), "payments".to_string()];
6971 let mut buf = BytesMut::new();
6972 encode_metadata_request(&mut buf, 12, Some(&topics), true).unwrap();
6973 let (got, allow, include_topic, include_cluster) =
6974 decode_metadata_request(&mut &buf[..], 12).unwrap();
6975 assert_eq!(got.as_deref(), Some(topics.as_slice()));
6976 assert!(allow);
6977 assert!(
6978 !include_topic,
6979 "encode_metadata_request must leave IncludeTopicAuthorizedOperations unset"
6980 );
6981 assert!(
6982 !include_cluster,
6983 "encode_metadata_request must leave IncludeClusterAuthorizedOperations unset"
6984 );
6985
6986 let mut all = BytesMut::new();
6987 encode_metadata_request(&mut all, 12, None, false).unwrap();
6988 let (got, allow, include_topic, include_cluster) =
6989 decode_metadata_request(&mut &all[..], 12).unwrap();
6990 assert!(got.is_none());
6991 assert!(!allow);
6992 assert!(!include_topic);
6993 assert!(!include_cluster);
6994
6995 let mut with = BytesMut::new();
6996 encode_metadata_request_with(&mut with, 12, Some(&topics), false, true).unwrap();
6997 let (got, allow, include_topic, include_cluster) =
6998 decode_metadata_request(&mut &with[..], 12).unwrap();
6999 assert_eq!(got.as_deref(), Some(topics.as_slice()));
7000 assert!(!allow);
7001 assert!(include_topic);
7002 assert!(!include_cluster);
7003 assert_ne!(
7004 &buf[..],
7005 &with[..],
7006 "IncludeTopicAuthorizedOperations true must not match the default request"
7007 );
7008 }
7009
7010 #[test]
7011 fn metadata_request_include_cluster_authorized_operations_matches_java() {
7012 let empty: [MetadataRequestTopic; 0] = [];
7021 for version in [8_i16, 9, 10] {
7022 let mut buf = BytesMut::new();
7023 encode_metadata_request_topics_with_include_cluster_authorized_operations(
7024 &mut buf,
7025 version,
7026 Some(&empty),
7027 true,
7028 false,
7029 true,
7030 )
7031 .unwrap();
7032 let mut cur = buf.as_ref();
7033 let (got, allow, include_topic, include_cluster) =
7034 decode_metadata_request_topics(&mut cur, version).unwrap();
7035 assert_eq!(got, Some(Vec::new()));
7036 assert!(allow);
7037 assert!(!include_topic);
7038 assert!(include_cluster);
7039 assert!(
7040 cur.is_empty(),
7041 "Metadata v{version} IncludeClusterAuthorizedOperations leftover-empty"
7042 );
7043 }
7044
7045 for version in [1_i16, 2, 3, 4, 5, 6, 7, 11, 12, 13] {
7046 let mut buf = BytesMut::new();
7047 encode_metadata_request_topics_with_include_cluster_authorized_operations(
7048 &mut buf,
7049 version,
7050 Some(&empty),
7051 true,
7052 false,
7053 true,
7054 )
7055 .unwrap();
7056 let mut cur = buf.as_ref();
7057 let (.., include_cluster) = decode_metadata_request_topics(&mut cur, version).unwrap();
7058 assert!(
7059 !include_cluster,
7060 "Metadata v{version} omits IncludeClusterAuthorizedOperations even when true"
7061 );
7062 assert!(
7063 cur.is_empty(),
7064 "Metadata v{version} IncludeClusterAuthorizedOperations leftover-empty"
7065 );
7066 let mut omit = BytesMut::new();
7067 encode_metadata_request_topics(&mut omit, version, Some(&empty), true, false).unwrap();
7068 assert_eq!(
7069 &buf[..],
7070 &omit[..],
7071 "Metadata v{version} encode omits IncludeClusterAuthorizedOperations even when true"
7072 );
7073 }
7074
7075 let mut with = BytesMut::new();
7076 encode_metadata_request_topics_with_include_cluster_authorized_operations(
7077 &mut with,
7078 8,
7079 Some(&empty),
7080 true,
7081 false,
7082 true,
7083 )
7084 .unwrap();
7085 let mut none = BytesMut::new();
7086 encode_metadata_request_topics(&mut none, 8, Some(&empty), true, false).unwrap();
7087 assert_ne!(
7088 &with[..],
7089 &none[..],
7090 "v8 IncludeClusterAuthorizedOperations is not always false"
7091 );
7092 assert_eq!(
7093 with.get(5..6),
7094 Some([1].as_slice()),
7095 "v8 classic IncludeClusterAuthorizedOperations follows empty Topics and AllowAutoTopicCreation"
7096 );
7097 assert_eq!(
7098 none.get(5..6),
7099 Some([0].as_slice()),
7100 "encode_metadata_request_topics still writes false IncludeClusterAuthorizedOperations"
7101 );
7102
7103 let mut v7_with = BytesMut::new();
7104 encode_metadata_request_topics_with_include_cluster_authorized_operations(
7105 &mut v7_with,
7106 7,
7107 Some(&empty),
7108 true,
7109 false,
7110 true,
7111 )
7112 .unwrap();
7113 assert_ne!(
7114 &v7_with[..],
7115 &with[..],
7116 "v8 adds IncludeClusterAuthorizedOperations after AllowAutoTopicCreation"
7117 );
7118
7119 let mut v9_with = BytesMut::new();
7120 encode_metadata_request_topics_with_include_cluster_authorized_operations(
7121 &mut v9_with,
7122 9,
7123 Some(&empty),
7124 true,
7125 false,
7126 true,
7127 )
7128 .unwrap();
7129 assert_ne!(
7130 &with[..],
7131 &v9_with[..],
7132 "v9 adds compact arrays/strings and tagged fields"
7133 );
7134 let mut v10_with = BytesMut::new();
7135 encode_metadata_request_topics_with_include_cluster_authorized_operations(
7136 &mut v10_with,
7137 10,
7138 Some(&empty),
7139 true,
7140 false,
7141 true,
7142 )
7143 .unwrap();
7144 assert_eq!(
7145 &v9_with[..],
7146 &v10_with[..],
7147 "empty-Topics IncludeClusterAuthorizedOperations bodies: v9 == v10"
7148 );
7149 let mut v11_with = BytesMut::new();
7150 encode_metadata_request_topics_with_include_cluster_authorized_operations(
7151 &mut v11_with,
7152 11,
7153 Some(&empty),
7154 true,
7155 false,
7156 true,
7157 )
7158 .unwrap();
7159 assert_ne!(
7160 &v10_with[..],
7161 &v11_with[..],
7162 "v11 drops IncludeClusterAuthorizedOperations"
7163 );
7164 }
7165
7166 #[test]
7167 fn metadata_v12_topic_id_request_is_compact() {
7168 const V12_ID: &[u8] = &[
7172 0x02, 0x74, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
7173 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
7174 ];
7175 let mut id = [0u8; 16];
7176 id[0] = b't';
7177 let topics = [MetadataRequestTopic::by_id(id)];
7178 let mut buf = BytesMut::new();
7179 encode_metadata_request_topics(&mut buf, 12, Some(&topics), false, false).unwrap();
7180 assert_eq!(&buf[..], V12_ID);
7181 let mut cur = &buf[..];
7182 let (got, allow, include_topic, include_cluster) =
7183 decode_metadata_request_topics(&mut cur, 12).unwrap();
7184 leftover_empty(&cur, "Metadata v12 TopicId request leftover").unwrap();
7185 let got = got.expect("Topics array");
7186 assert_eq!(got.as_slice(), topics.as_slice());
7187 assert!(!allow);
7188 assert!(!include_topic);
7189 assert!(!include_cluster);
7190 let mut cur = &buf[..];
7191 let (names_only, ..) = decode_metadata_request(&mut cur, 12).unwrap();
7192 leftover_empty(&cur, "Metadata v12 TopicId names-only leftover").unwrap();
7193 assert_eq!(
7194 names_only.as_deref(),
7195 Some(&[][..]),
7196 "name-only decode skips null-Name TopicId describes"
7197 );
7198 }
7199
7200 #[test]
7201 fn metadata_builder_matches_java() {
7202 let names = ["t".to_string()];
7203 let err = encode_metadata_request(&mut BytesMut::new(), 0, Some(&names), true).unwrap_err();
7204 assert!(
7205 matches!(err, Error::Unsupported(_)),
7206 "v0 is Java UnsupportedVersionException, got {err}"
7207 );
7208 assert!(err.to_string().contains("older than 1"), "got {err}");
7209 encode_metadata_request(&mut BytesMut::new(), 3, Some(&names), true).unwrap();
7210 let err =
7211 encode_metadata_request(&mut BytesMut::new(), 3, Some(&names), false).unwrap_err();
7212 assert!(
7213 matches!(err, Error::Unsupported(_)),
7214 "allowAutoTopicCreation false below v4 is Java UnsupportedVersionException, got {err}"
7215 );
7216 assert!(
7217 err.to_string().contains("allowAutoTopicCreation"),
7218 "got {err}"
7219 );
7220 encode_metadata_request(&mut BytesMut::new(), 4, Some(&names), false).unwrap();
7221
7222 let mut id = [0u8; 16];
7223 id[0] = 1;
7224 let by_id = [MetadataRequestTopic::by_id(id)];
7225 let err =
7226 encode_metadata_request_topics(&mut BytesMut::new(), 11, Some(&by_id), false, false)
7227 .unwrap_err();
7228 assert!(
7229 matches!(err, Error::Unsupported(_)),
7230 "null Name below v12 is Java UnsupportedVersionException, got {err}"
7231 );
7232 assert!(err.to_string().contains("null topic names"), "got {err}");
7233 encode_metadata_request_topics(&mut BytesMut::new(), 12, Some(&by_id), false, false)
7234 .unwrap();
7235 assert_eq!(
7236 MetadataRequestTopic::convert_from_names(["t"]),
7237 vec![MetadataRequestTopic::by_name("t")]
7238 );
7239 assert_eq!(
7240 MetadataRequestTopic::convert_from_ids([id]),
7241 vec![MetadataRequestTopic::by_id(id)]
7242 );
7243 assert!(MetadataRequest::is_all_topics(12, None));
7244 assert!(MetadataRequest::is_all_topics(0, Some(&[])));
7245 assert!(!MetadataRequest::is_all_topics(1, Some(&[])));
7246 let named = [MetadataRequestTopic::by_name("t")];
7247 assert!(!MetadataRequest::is_all_topics(12, Some(&named)));
7248 assert!(MetadataRequest::topic_ids(12, None).is_empty());
7249 assert!(MetadataRequest::topic_ids(0, Some(&[])).is_empty());
7250 assert!(MetadataRequest::topic_ids(9, Some(&named)).is_empty());
7251 assert_eq!(
7252 MetadataRequest::topic_ids(10, Some(&named)),
7253 vec![[0u8; 16]]
7254 );
7255 assert_eq!(MetadataRequest::topic_ids(12, Some(&by_id)), vec![id]);
7256 assert_eq!(MetadataRequest::topics(12, None), None);
7257 assert_eq!(MetadataRequest::topics(0, Some(&[])), None);
7258 assert_eq!(MetadataRequest::topics(1, Some(&[])), Some(Vec::new()));
7259 assert_eq!(
7260 MetadataRequest::topics(12, Some(&named)),
7261 Some(vec![Some("t")])
7262 );
7263 assert_eq!(MetadataRequest::topics(12, Some(&by_id)), Some(vec![None]));
7264 let named_err = MetadataRequestTopic::by_name("t")
7265 .error_result(crate::error::UNKNOWN_TOPIC_OR_PARTITION);
7266 assert_eq!(
7267 named_err,
7268 TopicMetadata::error(crate::error::UNKNOWN_TOPIC_OR_PARTITION, Some("t"), [0; 16])
7269 );
7270 assert!(!named_err.is_internal);
7271 assert!(named_err.partitions.is_empty());
7272 assert_eq!(
7273 named_err.topic_authorized_operations,
7274 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
7275 );
7276 let id_err =
7277 MetadataRequestTopic::by_id(id).error_result(crate::error::UNKNOWN_TOPIC_OR_PARTITION);
7278 assert_eq!(
7279 id_err,
7280 TopicMetadata::error(crate::error::UNKNOWN_TOPIC_OR_PARTITION, None, id)
7281 );
7282 assert_eq!(id_err.name.as_deref(), Some(""));
7283 let resp = MetadataResponse {
7284 throttle_time_ms: 0,
7285 brokers: Vec::new(),
7286 cluster_id: None,
7287 controller_id: MetadataResponse::NO_CONTROLLER_ID,
7288 topics: vec![named_err, id_err],
7289 cluster_authorized_operations: MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
7290 error_code: crate::error::UNKNOWN_TOPIC_OR_PARTITION,
7291 };
7292 let mut buf = BytesMut::new();
7293 encode_metadata_response(&mut buf, 13, &resp).unwrap();
7294 let mut cur = buf.as_ref();
7295 let decoded = decode_metadata_response(&mut cur, 13).unwrap();
7296 leftover_empty(&cur, "Metadata getErrorResponse v13").unwrap();
7297 assert_eq!(decoded, resp);
7298 let named_id = [MetadataRequestTopic {
7299 name: Some("t".into()),
7300 topic_id: id,
7301 }];
7302 let err =
7303 encode_metadata_request_topics(&mut BytesMut::new(), 11, Some(&named_id), true, false)
7304 .unwrap_err();
7305 assert!(err.to_string().contains("non-zero topic IDs"), "got {err}");
7306 }
7307
7308 #[test]
7309 fn metadata_request_error_response_matches_java() {
7310 let code = crate::error::UNKNOWN_TOPIC_OR_PARTITION;
7318 let none = MetadataRequest::error_response(None, code);
7319 assert_eq!(none.error_code, code);
7320 assert_eq!(none.throttle_time_ms, 0);
7321 assert_eq!(
7322 none.cluster_authorized_operations(),
7323 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
7324 );
7325 assert!(none.brokers.is_empty());
7326 assert!(none.cluster_id.is_none());
7327 assert_eq!(none.controller_id, MetadataResponse::NO_CONTROLLER_ID);
7328 assert!(none.topics.is_empty());
7329
7330 let empty_list = MetadataRequest::error_response(Some(&[]), code);
7331 assert_eq!(empty_list, none);
7332
7333 let mut id = [0u8; 16];
7334 id[0] = 1;
7335 let named = MetadataRequestTopic::by_name("orders");
7336 let by_id = MetadataRequestTopic::by_id(id);
7337 let dup = MetadataRequestTopic::by_name("orders");
7338 let topics = [named.clone(), by_id.clone(), dup];
7339 let grouped = MetadataRequest::error_response(Some(topics.as_slice()), code);
7340 assert_eq!(grouped.error_code, code);
7341 assert!(grouped.brokers.is_empty());
7342 assert_eq!(grouped.topics.len(), 3);
7343 assert_eq!(grouped.topics[0], named.error_result(code));
7344 assert_eq!(grouped.topics[1].name.as_deref(), Some(""));
7345 assert_eq!(grouped.topics[1].topic_id, id);
7346 assert_eq!(grouped.topics[2].name.as_deref(), Some("orders"));
7347 leftover_metadata_error_response(1, Some(std::slice::from_ref(&named)));
7348 leftover_metadata_error_response(1, None);
7349 leftover_metadata_error_response(9, Some(std::slice::from_ref(&named)));
7350 leftover_metadata_error_response(9, Some(&[]));
7351 leftover_metadata_error_response(13, Some(topics.as_slice()));
7352 leftover_metadata_error_response(13, None);
7353 }
7354
7355 #[test]
7356 fn metadata_response_cluster_authorized_operations_matches_java() {
7357 fn empty(ops: i32) -> MetadataResponse {
7368 MetadataResponse {
7369 throttle_time_ms: 0,
7370 brokers: Vec::new(),
7371 cluster_id: None,
7372 controller_id: MetadataResponse::NO_CONTROLLER_ID,
7373 topics: Vec::new(),
7374 cluster_authorized_operations: ops,
7375 error_code: 0,
7376 }
7377 }
7378 let with = empty(1);
7379 let omitted = empty(MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED);
7380 assert_eq!(with.cluster_authorized_operations(), 1);
7381 assert_eq!(
7382 omitted.cluster_authorized_operations(),
7383 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
7384 );
7385 assert_eq!(
7386 MetadataRequest::error_response(None, 0).cluster_authorized_operations(),
7387 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
7388 "Java MetadataRequest.getErrorResponse does not set ClusterAuthorizedOperations"
7389 );
7390
7391 for version in [8_i16, 9, 10] {
7392 let mut buf = BytesMut::new();
7393 encode_metadata_response(&mut buf, version, &with).unwrap();
7394 let mut cur = buf.as_ref();
7395 let decoded = decode_metadata_response(&mut cur, version).unwrap();
7396 assert_eq!(decoded.cluster_authorized_operations(), 1);
7397 assert_eq!(decoded, with);
7398 assert!(
7399 cur.is_empty(),
7400 "Metadata v{version} ClusterAuthorizedOperations leftover-empty"
7401 );
7402 }
7403
7404 for version in [1_i16, 2, 3, 4, 5, 6, 7, 11, 12, 13] {
7405 let mut buf = BytesMut::new();
7406 encode_metadata_response(&mut buf, version, &with).unwrap();
7407 let mut cur = buf.as_ref();
7408 let decoded = decode_metadata_response(&mut cur, version).unwrap();
7409 assert_eq!(decoded, omitted);
7410 assert_eq!(
7411 decoded.cluster_authorized_operations(),
7412 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED,
7413 "Metadata v{version} omits ClusterAuthorizedOperations even when the body is non-default"
7414 );
7415 assert!(
7416 cur.is_empty(),
7417 "Metadata v{version} ClusterAuthorizedOperations leftover-empty"
7418 );
7419 let mut omit_buf = BytesMut::new();
7420 encode_metadata_response(&mut omit_buf, version, &omitted).unwrap();
7421 assert_eq!(
7422 &buf[..],
7423 &omit_buf[..],
7424 "Metadata v{version} encode omits ClusterAuthorizedOperations even when the body is non-default"
7425 );
7426 }
7427
7428 let mut v8_with = BytesMut::new();
7429 encode_metadata_response(&mut v8_with, 8, &with).unwrap();
7430 let mut v8_omit = BytesMut::new();
7431 encode_metadata_response(&mut v8_omit, 8, &omitted).unwrap();
7432 assert_ne!(
7433 &v8_with[..],
7434 &v8_omit[..],
7435 "v8 ClusterAuthorizedOperations is not always the JSON default AUTHORIZED_OPERATIONS_OMITTED"
7436 );
7437 assert_eq!(
7438 v8_with.get(18..22),
7439 Some([0, 0, 0, 1].as_slice()),
7440 "v8 classic ClusterAuthorizedOperations follows empty Topics"
7441 );
7442 assert_eq!(
7443 v8_omit.get(18..22),
7444 Some([0x80, 0, 0, 0].as_slice()),
7445 "v8 classic ClusterAuthorizedOperations JSON default is AUTHORIZED_OPERATIONS_OMITTED"
7446 );
7447
7448 let mut v7_with = BytesMut::new();
7449 encode_metadata_response(&mut v7_with, 7, &with).unwrap();
7450 assert_ne!(
7451 &v7_with[..],
7452 &v8_with[..],
7453 "v8 adds ClusterAuthorizedOperations after Topics"
7454 );
7455
7456 let mut v9_with = BytesMut::new();
7457 encode_metadata_response(&mut v9_with, 9, &with).unwrap();
7458 assert_ne!(
7459 &v8_with[..],
7460 &v9_with[..],
7461 "v9 adds compact arrays/strings and tagged fields"
7462 );
7463 let mut v10_with = BytesMut::new();
7464 encode_metadata_response(&mut v10_with, 10, &with).unwrap();
7465 assert_eq!(
7466 &v9_with[..],
7467 &v10_with[..],
7468 "empty-Topics ClusterAuthorizedOperations bodies: v9 == v10"
7469 );
7470 let mut v11_with = BytesMut::new();
7471 encode_metadata_response(&mut v11_with, 11, &with).unwrap();
7472 assert_ne!(
7473 &v10_with[..],
7474 &v11_with[..],
7475 "v11 drops ClusterAuthorizedOperations"
7476 );
7477 }
7478
7479 fn leftover_metadata_error_response(version: i16, topics: Option<&[MetadataRequestTopic]>) {
7480 let resp =
7481 MetadataRequest::error_response(topics, crate::error::UNKNOWN_TOPIC_OR_PARTITION);
7482 let mut buf = BytesMut::new();
7483 encode_metadata_response(&mut buf, version, &resp).unwrap();
7484 let mut cur = buf.as_ref();
7485 let decoded = decode_metadata_response(&mut cur, version).unwrap();
7486 leftover_empty(
7487 &cur,
7488 match (version, topics.filter(|t| !t.is_empty()).is_none()) {
7489 (1, false) => "Metadata v1 Request.getErrorResponse leftover-empty",
7490 (1, true) => "Metadata v1 Request.getErrorResponse empty leftover-empty",
7491 (9, false) => "Metadata v9 Request.getErrorResponse leftover-empty",
7492 (9, true) => "Metadata v9 Request.getErrorResponse empty leftover-empty",
7493 (13, false) => "Metadata v13 Request.getErrorResponse leftover-empty",
7494 _ => "Metadata v13 Request.getErrorResponse empty leftover-empty",
7495 },
7496 )
7497 .unwrap();
7498 assert_eq!(decoded.topics, resp.topics);
7499 assert!(decoded.brokers.is_empty());
7500 assert_eq!(decoded.controller_id, MetadataResponse::NO_CONTROLLER_ID);
7501 assert_eq!(
7502 decoded.cluster_authorized_operations(),
7503 MetadataResponse::AUTHORIZED_OPERATIONS_OMITTED
7504 );
7505 if version >= 13 {
7506 assert_eq!(decoded.error_code, resp.error_code);
7507 } else {
7508 assert_eq!(decoded.error_code, 0);
7509 }
7510 }
7511}