Skip to main content

kacrab_protocol/generated/
produce_response.rs

1//! Generated from ProduceResponse.json - DO NOT EDIT
2#![allow(
3    missing_docs,
4    clippy::all,
5    clippy::pedantic,
6    clippy::nursery,
7    clippy::arithmetic_side_effects,
8    reason = "Generated protocol modules mirror Kafka's schema shape and intentionally trade \
9              hand-written lint style for reproducible wire-code output."
10)]
11use bytes::{Bytes, BytesMut};
12
13use crate::*;
14
15#[derive(Debug, Clone, PartialEq)]
16pub struct ProduceResponseData {
17    /// Each produce response.
18    pub responses: Vec<TopicProduceResponse>,
19    /// The duration in milliseconds for which the request was throttled due to a quota violation,
20    /// or zero if the request did not violate any quota.
21    pub throttle_time_ms: i32,
22    /// Endpoints for all current-leaders enumerated in PartitionProduceResponses, with errors
23    /// NOT_LEADER_OR_FOLLOWER.
24    pub node_endpoints: Vec<NodeEndpoint>,
25    pub _unknown_tagged_fields: Vec<RawTaggedField>,
26}
27impl Default for ProduceResponseData {
28    fn default() -> Self {
29        Self {
30            responses: Vec::new(),
31            throttle_time_ms: 0i32,
32            node_endpoints: Vec::new(),
33            _unknown_tagged_fields: Vec::new(),
34        }
35    }
36}
37impl ProduceResponseData {
38    pub fn with_responses(mut self, value: Vec<TopicProduceResponse>) -> Self {
39        self.responses = value;
40        self
41    }
42    pub fn with_throttle_time_ms(mut self, value: i32) -> Self {
43        self.throttle_time_ms = value;
44        self
45    }
46    pub fn with_node_endpoints(mut self, value: Vec<NodeEndpoint>) -> Self {
47        self.node_endpoints = value;
48        self
49    }
50    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
51        if version < 3 || version > 13 {
52            return Err(UnsupportedVersion::new(0, version).into());
53        }
54        let responses;
55        let throttle_time_ms;
56        let mut node_endpoints = Vec::new();
57        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
58        if version >= 9 {
59            responses = {
60                let len = read_compact_array_length(buf)?;
61                let mut arr = Vec::with_capacity(len.max(0) as usize);
62                for _ in 0..len {
63                    arr.push(TopicProduceResponse::read(buf, version)?);
64                }
65                arr
66            };
67        } else {
68            responses = {
69                let len = read_array_length(buf)?;
70                let mut arr = Vec::with_capacity(len.max(0) as usize);
71                for _ in 0..len {
72                    arr.push(TopicProduceResponse::read(buf, version)?);
73                }
74                arr
75            };
76        }
77        throttle_time_ms = read_i32(buf)?;
78        if version >= 9 {
79            let tagged_fields = read_tagged_fields(buf)?;
80            for field in &tagged_fields {
81                match field.tag {
82                    0 => {
83                        if version >= 10 {
84                            let mut tag_buf = field.data.clone();
85                            node_endpoints = {
86                                let len = read_compact_array_length(&mut tag_buf)?;
87                                let mut arr = Vec::with_capacity(len.max(0) as usize);
88                                for _ in 0..len {
89                                    arr.push(NodeEndpoint::read(&mut tag_buf, version)?);
90                                }
91                                arr
92                            };
93                        }
94                    },
95                    _ => {
96                        _unknown_tagged_fields.push(field.clone());
97                    },
98                }
99            }
100        }
101        Ok(Self {
102            responses,
103            throttle_time_ms,
104            node_endpoints,
105            _unknown_tagged_fields,
106        })
107    }
108    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
109        if version < 3 || version > 13 {
110            return Err(UnsupportedVersion::new(0, version).into());
111        }
112        if version >= 9 {
113            write_compact_array_length(buf, self.responses.len() as i32);
114            for el in &self.responses {
115                el.write(buf, version)?;
116            }
117        } else {
118            write_array_length(buf, self.responses.len() as i32);
119            for el in &self.responses {
120                el.write(buf, version)?;
121            }
122        }
123        write_i32(buf, self.throttle_time_ms);
124        if version >= 9 {
125            let mut known_tagged_fields: Vec<RawTaggedField> = Vec::new();
126            if version >= 10 && !self.node_endpoints.is_empty() {
127                let mut tag_buf = BytesMut::new();
128                write_compact_array_length(&mut tag_buf, self.node_endpoints.len() as i32);
129                for el in &self.node_endpoints {
130                    el.write(&mut tag_buf, version)?;
131                }
132                known_tagged_fields.push(RawTaggedField {
133                    tag: 0,
134                    data: tag_buf.freeze(),
135                });
136            }
137            let mut all_tags = known_tagged_fields;
138            all_tags.extend(self._unknown_tagged_fields.iter().cloned());
139            all_tags.sort_by_key(|f| f.tag);
140            write_tagged_fields(buf, &all_tags)?;
141        }
142        Ok(())
143    }
144    pub fn encoded_len(&self, version: i16) -> Result<usize> {
145        if version < 3 || version > 13 {
146            return Err(UnsupportedVersion::new(0, version).into());
147        }
148        let mut len: usize = 0;
149        if version >= 9 {
150            len += compact_array_length_len(self.responses.len() as i32);
151            for el in &self.responses {
152                len += el.encoded_len(version)?;
153            }
154        } else {
155            len += array_length_len();
156            for el in &self.responses {
157                len += el.encoded_len(version)?;
158            }
159        }
160        len += 4;
161        if version >= 9 {
162            let mut known_tagged_fields: Vec<RawTaggedField> = Vec::new();
163            if version >= 10 && !self.node_endpoints.is_empty() {
164                let mut tag_buf = BytesMut::new();
165                write_compact_array_length(&mut tag_buf, self.node_endpoints.len() as i32);
166                for el in &self.node_endpoints {
167                    el.write(&mut tag_buf, version)?;
168                }
169                known_tagged_fields.push(RawTaggedField {
170                    tag: 0,
171                    data: tag_buf.freeze(),
172                });
173            }
174            let mut all_tags = known_tagged_fields;
175            all_tags.extend(self._unknown_tagged_fields.iter().cloned());
176            all_tags.sort_by_key(|f| f.tag);
177            len += tagged_fields_len(&all_tags)?;
178        }
179        Ok(len)
180    }
181}
182#[derive(Debug, Clone, PartialEq)]
183pub struct TopicProduceResponse {
184    /// The topic name.
185    pub name: KafkaString,
186    /// The unique topic ID
187    pub topic_id: KafkaUuid,
188    /// Each partition that we produced to within the topic.
189    pub partition_responses: Vec<PartitionProduceResponse>,
190    pub _unknown_tagged_fields: Vec<RawTaggedField>,
191}
192impl Default for TopicProduceResponse {
193    fn default() -> Self {
194        Self {
195            name: KafkaString::default(),
196            topic_id: KafkaUuid::ZERO,
197            partition_responses: Vec::new(),
198            _unknown_tagged_fields: Vec::new(),
199        }
200    }
201}
202impl TopicProduceResponse {
203    pub fn with_name(mut self, value: KafkaString) -> Self {
204        self.name = value;
205        self
206    }
207    pub fn with_topic_id(mut self, value: KafkaUuid) -> Self {
208        self.topic_id = value;
209        self
210    }
211    pub fn with_partition_responses(mut self, value: Vec<PartitionProduceResponse>) -> Self {
212        self.partition_responses = value;
213        self
214    }
215    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
216        let mut name = KafkaString::default();
217        let mut topic_id = KafkaUuid::ZERO;
218        let partition_responses;
219        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
220        if version <= 12 {
221            if version >= 9 {
222                name = read_compact_string(buf)?;
223            } else {
224                name = read_string(buf)?;
225            }
226        }
227        if version >= 13 {
228            topic_id = read_uuid(buf)?;
229        }
230        if version >= 9 {
231            partition_responses = {
232                let len = read_compact_array_length(buf)?;
233                let mut arr = Vec::with_capacity(len.max(0) as usize);
234                for _ in 0..len {
235                    arr.push(PartitionProduceResponse::read(buf, version)?);
236                }
237                arr
238            };
239        } else {
240            partition_responses = {
241                let len = read_array_length(buf)?;
242                let mut arr = Vec::with_capacity(len.max(0) as usize);
243                for _ in 0..len {
244                    arr.push(PartitionProduceResponse::read(buf, version)?);
245                }
246                arr
247            };
248        }
249        if version >= 9 {
250            let tagged_fields = read_tagged_fields(buf)?;
251            for field in &tagged_fields {
252                match field.tag {
253                    _ => {
254                        _unknown_tagged_fields.push(field.clone());
255                    },
256                }
257            }
258        }
259        Ok(Self {
260            name,
261            topic_id,
262            partition_responses,
263            _unknown_tagged_fields,
264        })
265    }
266    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
267        if version <= 12 {
268            if version >= 9 {
269                write_compact_string(buf, &self.name)?;
270            } else {
271                write_string(buf, &self.name)?;
272            }
273        } else if self.name != KafkaString::default() {
274            return Err(UnsupportedFieldVersion::new(0, "name", version).into());
275        }
276        if version >= 13 {
277            write_uuid(buf, &self.topic_id);
278        } else if self.topic_id != KafkaUuid::ZERO {
279            return Err(UnsupportedFieldVersion::new(0, "topic_id", version).into());
280        }
281        if version >= 9 {
282            write_compact_array_length(buf, self.partition_responses.len() as i32);
283            for el in &self.partition_responses {
284                el.write(buf, version)?;
285            }
286        } else {
287            write_array_length(buf, self.partition_responses.len() as i32);
288            for el in &self.partition_responses {
289                el.write(buf, version)?;
290            }
291        }
292        if version >= 9 {
293            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
294            all_tags.sort_by_key(|f| f.tag);
295            write_tagged_fields(buf, &all_tags)?;
296        }
297        Ok(())
298    }
299    pub fn encoded_len(&self, version: i16) -> Result<usize> {
300        let mut len: usize = 0;
301        if version <= 12 {
302            if version >= 9 {
303                len += compact_string_len(&self.name)?;
304            } else {
305                len += string_len(&self.name)?;
306            }
307        } else if self.name != KafkaString::default() {
308            return Err(UnsupportedFieldVersion::new(0, "name", version).into());
309        }
310        if version >= 13 {
311            len += 16;
312        } else if self.topic_id != KafkaUuid::ZERO {
313            return Err(UnsupportedFieldVersion::new(0, "topic_id", version).into());
314        }
315        if version >= 9 {
316            len += compact_array_length_len(self.partition_responses.len() as i32);
317            for el in &self.partition_responses {
318                len += el.encoded_len(version)?;
319            }
320        } else {
321            len += array_length_len();
322            for el in &self.partition_responses {
323                len += el.encoded_len(version)?;
324            }
325        }
326        if version >= 9 {
327            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
328            all_tags.sort_by_key(|f| f.tag);
329            len += tagged_fields_len(&all_tags)?;
330        }
331        Ok(len)
332    }
333}
334#[derive(Debug, Clone, PartialEq)]
335pub struct PartitionProduceResponse {
336    /// The partition index.
337    pub index: i32,
338    /// The error code, or 0 if there was no error.
339    pub error_code: i16,
340    /// The base offset.
341    pub base_offset: i64,
342    /// The timestamp returned by broker after appending the messages. If CreateTime is used for
343    /// the topic, the timestamp will be -1.  If LogAppendTime is used for the topic, the timestamp
344    /// will be the broker local time when the messages are appended.
345    pub log_append_time_ms: i64,
346    /// The log start offset.
347    pub log_start_offset: i64,
348    /// The batch indices of records that caused the batch to be dropped.
349    pub record_errors: Vec<BatchIndexAndErrorMessage>,
350    /// The global error message summarizing the common root cause of the records that caused the
351    /// batch to be dropped.
352    pub error_message: Option<KafkaString>,
353    /// The leader broker that the producer should use for future requests.
354    pub current_leader: LeaderIdAndEpoch,
355    pub _unknown_tagged_fields: Vec<RawTaggedField>,
356}
357impl Default for PartitionProduceResponse {
358    fn default() -> Self {
359        Self {
360            index: 0_i32,
361            error_code: 0_i16,
362            base_offset: 0_i64,
363            log_append_time_ms: -1i64,
364            log_start_offset: -1i64,
365            record_errors: Vec::new(),
366            error_message: None,
367            current_leader: LeaderIdAndEpoch::default(),
368            _unknown_tagged_fields: Vec::new(),
369        }
370    }
371}
372impl PartitionProduceResponse {
373    pub fn with_index(mut self, value: i32) -> Self {
374        self.index = value;
375        self
376    }
377    pub fn with_error_code(mut self, value: i16) -> Self {
378        self.error_code = value;
379        self
380    }
381    pub fn with_base_offset(mut self, value: i64) -> Self {
382        self.base_offset = value;
383        self
384    }
385    pub fn with_log_append_time_ms(mut self, value: i64) -> Self {
386        self.log_append_time_ms = value;
387        self
388    }
389    pub fn with_log_start_offset(mut self, value: i64) -> Self {
390        self.log_start_offset = value;
391        self
392    }
393    pub fn with_record_errors(mut self, value: Vec<BatchIndexAndErrorMessage>) -> Self {
394        self.record_errors = value;
395        self
396    }
397    pub fn with_error_message(mut self, value: Option<KafkaString>) -> Self {
398        self.error_message = value;
399        self
400    }
401    pub fn with_current_leader(mut self, value: LeaderIdAndEpoch) -> Self {
402        self.current_leader = value;
403        self
404    }
405    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
406        let index;
407        let error_code;
408        let base_offset;
409        let log_append_time_ms;
410        let mut log_start_offset = -1i64;
411        let mut record_errors = Vec::new();
412        let mut error_message = None;
413        let mut current_leader = LeaderIdAndEpoch::default();
414        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
415        index = read_i32(buf)?;
416        error_code = read_i16(buf)?;
417        base_offset = read_i64(buf)?;
418        log_append_time_ms = read_i64(buf)?;
419        if version >= 5 {
420            log_start_offset = read_i64(buf)?;
421        }
422        if version >= 8 {
423            if version >= 9 {
424                record_errors = {
425                    let len = read_compact_array_length(buf)?;
426                    let mut arr = Vec::with_capacity(len.max(0) as usize);
427                    for _ in 0..len {
428                        arr.push(BatchIndexAndErrorMessage::read(buf, version)?);
429                    }
430                    arr
431                };
432            } else {
433                record_errors = {
434                    let len = read_array_length(buf)?;
435                    let mut arr = Vec::with_capacity(len.max(0) as usize);
436                    for _ in 0..len {
437                        arr.push(BatchIndexAndErrorMessage::read(buf, version)?);
438                    }
439                    arr
440                };
441            }
442        }
443        if version >= 8 {
444            if version >= 9 {
445                error_message = read_compact_nullable_string(buf)?;
446            } else {
447                error_message = read_nullable_string(buf)?;
448            }
449        }
450        if version >= 9 {
451            let tagged_fields = read_tagged_fields(buf)?;
452            for field in &tagged_fields {
453                match field.tag {
454                    0 => {
455                        if version >= 10 {
456                            let mut tag_buf = field.data.clone();
457                            current_leader = LeaderIdAndEpoch::read(&mut tag_buf, version)?;
458                        }
459                    },
460                    _ => {
461                        _unknown_tagged_fields.push(field.clone());
462                    },
463                }
464            }
465        }
466        Ok(Self {
467            index,
468            error_code,
469            base_offset,
470            log_append_time_ms,
471            log_start_offset,
472            record_errors,
473            error_message,
474            current_leader,
475            _unknown_tagged_fields,
476        })
477    }
478    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
479        write_i32(buf, self.index);
480        write_i16(buf, self.error_code);
481        write_i64(buf, self.base_offset);
482        write_i64(buf, self.log_append_time_ms);
483        if version >= 5 {
484            write_i64(buf, self.log_start_offset);
485        } else if self.log_start_offset != -1i64 {
486            return Err(UnsupportedFieldVersion::new(0, "log_start_offset", version).into());
487        }
488        if version >= 8 {
489            if version >= 9 {
490                write_compact_array_length(buf, self.record_errors.len() as i32);
491                for el in &self.record_errors {
492                    el.write(buf, version)?;
493                }
494            } else {
495                write_array_length(buf, self.record_errors.len() as i32);
496                for el in &self.record_errors {
497                    el.write(buf, version)?;
498                }
499            }
500        } else if self.record_errors != Vec::new() {
501            return Err(UnsupportedFieldVersion::new(0, "record_errors", version).into());
502        }
503        if version >= 8 {
504            if version >= 9 {
505                write_compact_nullable_string(buf, self.error_message.as_ref())?;
506            } else {
507                write_nullable_string(buf, self.error_message.as_ref())?;
508            }
509        } else if self.error_message != None {
510            return Err(UnsupportedFieldVersion::new(0, "error_message", version).into());
511        }
512        if version >= 9 {
513            let mut known_tagged_fields: Vec<RawTaggedField> = Vec::new();
514            if version >= 10 && self.current_leader != LeaderIdAndEpoch::default() {
515                let mut tag_buf = BytesMut::new();
516                self.current_leader.write(&mut tag_buf, version)?;
517                known_tagged_fields.push(RawTaggedField {
518                    tag: 0,
519                    data: tag_buf.freeze(),
520                });
521            }
522            let mut all_tags = known_tagged_fields;
523            all_tags.extend(self._unknown_tagged_fields.iter().cloned());
524            all_tags.sort_by_key(|f| f.tag);
525            write_tagged_fields(buf, &all_tags)?;
526        }
527        Ok(())
528    }
529    pub fn encoded_len(&self, version: i16) -> Result<usize> {
530        let mut len: usize = 0;
531        len += 4;
532        len += 2;
533        len += 8;
534        len += 8;
535        if version >= 5 {
536            len += 8;
537        } else if self.log_start_offset != -1i64 {
538            return Err(UnsupportedFieldVersion::new(0, "log_start_offset", version).into());
539        }
540        if version >= 8 {
541            if version >= 9 {
542                len += compact_array_length_len(self.record_errors.len() as i32);
543                for el in &self.record_errors {
544                    len += el.encoded_len(version)?;
545                }
546            } else {
547                len += array_length_len();
548                for el in &self.record_errors {
549                    len += el.encoded_len(version)?;
550                }
551            }
552        } else if self.record_errors != Vec::new() {
553            return Err(UnsupportedFieldVersion::new(0, "record_errors", version).into());
554        }
555        if version >= 8 {
556            if version >= 9 {
557                len += compact_nullable_string_len(self.error_message.as_ref())?;
558            } else {
559                len += nullable_string_len(self.error_message.as_ref())?;
560            }
561        } else if self.error_message != None {
562            return Err(UnsupportedFieldVersion::new(0, "error_message", version).into());
563        }
564        if version >= 9 {
565            let mut known_tagged_fields: Vec<RawTaggedField> = Vec::new();
566            if version >= 10 && self.current_leader != LeaderIdAndEpoch::default() {
567                let mut tag_buf = BytesMut::new();
568                self.current_leader.write(&mut tag_buf, version)?;
569                known_tagged_fields.push(RawTaggedField {
570                    tag: 0,
571                    data: tag_buf.freeze(),
572                });
573            }
574            let mut all_tags = known_tagged_fields;
575            all_tags.extend(self._unknown_tagged_fields.iter().cloned());
576            all_tags.sort_by_key(|f| f.tag);
577            len += tagged_fields_len(&all_tags)?;
578        }
579        Ok(len)
580    }
581}
582#[derive(Debug, Clone, PartialEq)]
583pub struct BatchIndexAndErrorMessage {
584    /// The batch index of the record that caused the batch to be dropped.
585    pub batch_index: i32,
586    /// The error message of the record that caused the batch to be dropped.
587    pub batch_index_error_message: Option<KafkaString>,
588    pub _unknown_tagged_fields: Vec<RawTaggedField>,
589}
590impl Default for BatchIndexAndErrorMessage {
591    fn default() -> Self {
592        Self {
593            batch_index: 0_i32,
594            batch_index_error_message: None,
595            _unknown_tagged_fields: Vec::new(),
596        }
597    }
598}
599impl BatchIndexAndErrorMessage {
600    pub fn with_batch_index(mut self, value: i32) -> Self {
601        self.batch_index = value;
602        self
603    }
604    pub fn with_batch_index_error_message(mut self, value: Option<KafkaString>) -> Self {
605        self.batch_index_error_message = value;
606        self
607    }
608    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
609        let batch_index;
610        let batch_index_error_message;
611        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
612        batch_index = read_i32(buf)?;
613        if version >= 9 {
614            batch_index_error_message = read_compact_nullable_string(buf)?;
615        } else {
616            batch_index_error_message = read_nullable_string(buf)?;
617        }
618        if version >= 9 {
619            let tagged_fields = read_tagged_fields(buf)?;
620            for field in &tagged_fields {
621                match field.tag {
622                    _ => {
623                        _unknown_tagged_fields.push(field.clone());
624                    },
625                }
626            }
627        }
628        Ok(Self {
629            batch_index,
630            batch_index_error_message,
631            _unknown_tagged_fields,
632        })
633    }
634    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
635        write_i32(buf, self.batch_index);
636        if version >= 9 {
637            write_compact_nullable_string(buf, self.batch_index_error_message.as_ref())?;
638        } else {
639            write_nullable_string(buf, self.batch_index_error_message.as_ref())?;
640        }
641        if version >= 9 {
642            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
643            all_tags.sort_by_key(|f| f.tag);
644            write_tagged_fields(buf, &all_tags)?;
645        }
646        Ok(())
647    }
648    pub fn encoded_len(&self, version: i16) -> Result<usize> {
649        let mut len: usize = 0;
650        len += 4;
651        if version >= 9 {
652            len += compact_nullable_string_len(self.batch_index_error_message.as_ref())?;
653        } else {
654            len += nullable_string_len(self.batch_index_error_message.as_ref())?;
655        }
656        if version >= 9 {
657            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
658            all_tags.sort_by_key(|f| f.tag);
659            len += tagged_fields_len(&all_tags)?;
660        }
661        Ok(len)
662    }
663}
664#[derive(Debug, Clone, PartialEq)]
665pub struct LeaderIdAndEpoch {
666    /// The ID of the current leader or -1 if the leader is unknown.
667    pub leader_id: i32,
668    /// The latest known leader epoch.
669    pub leader_epoch: i32,
670    pub _unknown_tagged_fields: Vec<RawTaggedField>,
671}
672impl Default for LeaderIdAndEpoch {
673    fn default() -> Self {
674        Self {
675            leader_id: -1i32,
676            leader_epoch: -1i32,
677            _unknown_tagged_fields: Vec::new(),
678        }
679    }
680}
681impl LeaderIdAndEpoch {
682    pub fn with_leader_id(mut self, value: i32) -> Self {
683        self.leader_id = value;
684        self
685    }
686    pub fn with_leader_epoch(mut self, value: i32) -> Self {
687        self.leader_epoch = value;
688        self
689    }
690    pub fn read(buf: &mut Bytes, _version: i16) -> Result<Self> {
691        let leader_id;
692        let leader_epoch;
693        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
694        leader_id = read_i32(buf)?;
695        leader_epoch = read_i32(buf)?;
696        let tagged_fields = read_tagged_fields(buf)?;
697        for field in &tagged_fields {
698            match field.tag {
699                _ => {
700                    _unknown_tagged_fields.push(field.clone());
701                },
702            }
703        }
704        Ok(Self {
705            leader_id,
706            leader_epoch,
707            _unknown_tagged_fields,
708        })
709    }
710    pub fn write(&self, buf: &mut BytesMut, _version: i16) -> Result<()> {
711        write_i32(buf, self.leader_id);
712        write_i32(buf, self.leader_epoch);
713        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
714        all_tags.sort_by_key(|f| f.tag);
715        write_tagged_fields(buf, &all_tags)?;
716        Ok(())
717    }
718    pub fn encoded_len(&self, _version: i16) -> Result<usize> {
719        let mut len: usize = 0;
720        len += 4;
721        len += 4;
722        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
723        all_tags.sort_by_key(|f| f.tag);
724        len += tagged_fields_len(&all_tags)?;
725        Ok(len)
726    }
727}
728#[derive(Debug, Clone, PartialEq)]
729pub struct NodeEndpoint {
730    /// The ID of the associated node.
731    pub node_id: i32,
732    /// The node's hostname.
733    pub host: KafkaString,
734    /// The node's port.
735    pub port: i32,
736    /// The rack of the node, or null if it has not been assigned to a rack.
737    pub rack: Option<KafkaString>,
738    pub _unknown_tagged_fields: Vec<RawTaggedField>,
739}
740impl Default for NodeEndpoint {
741    fn default() -> Self {
742        Self {
743            node_id: 0_i32,
744            host: KafkaString::default(),
745            port: 0_i32,
746            rack: None,
747            _unknown_tagged_fields: Vec::new(),
748        }
749    }
750}
751impl NodeEndpoint {
752    pub fn with_node_id(mut self, value: i32) -> Self {
753        self.node_id = value;
754        self
755    }
756    pub fn with_host(mut self, value: KafkaString) -> Self {
757        self.host = value;
758        self
759    }
760    pub fn with_port(mut self, value: i32) -> Self {
761        self.port = value;
762        self
763    }
764    pub fn with_rack(mut self, value: Option<KafkaString>) -> Self {
765        self.rack = value;
766        self
767    }
768    pub fn read(buf: &mut Bytes, _version: i16) -> Result<Self> {
769        let node_id;
770        let host;
771        let port;
772        let rack;
773        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
774        node_id = read_i32(buf)?;
775        host = read_compact_string(buf)?;
776        port = read_i32(buf)?;
777        rack = read_compact_nullable_string(buf)?;
778        let tagged_fields = read_tagged_fields(buf)?;
779        for field in &tagged_fields {
780            match field.tag {
781                _ => {
782                    _unknown_tagged_fields.push(field.clone());
783                },
784            }
785        }
786        Ok(Self {
787            node_id,
788            host,
789            port,
790            rack,
791            _unknown_tagged_fields,
792        })
793    }
794    pub fn write(&self, buf: &mut BytesMut, _version: i16) -> Result<()> {
795        write_i32(buf, self.node_id);
796        write_compact_string(buf, &self.host)?;
797        write_i32(buf, self.port);
798        write_compact_nullable_string(buf, self.rack.as_ref())?;
799        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
800        all_tags.sort_by_key(|f| f.tag);
801        write_tagged_fields(buf, &all_tags)?;
802        Ok(())
803    }
804    pub fn encoded_len(&self, _version: i16) -> Result<usize> {
805        let mut len: usize = 0;
806        len += 4;
807        len += compact_string_len(&self.host)?;
808        len += 4;
809        len += compact_nullable_string_len(self.rack.as_ref())?;
810        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
811        all_tags.sort_by_key(|f| f.tag);
812        len += tagged_fields_len(&all_tags)?;
813        Ok(len)
814    }
815}