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