Skip to main content

kacrab_protocol/generated/
describe_quorum_response.rs

1//! Generated from DescribeQuorumResponse.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 DescribeQuorumResponseData {
17    /// The top level error code.
18    pub error_code: i16,
19    /// The error message, or null if there was no error.
20    pub error_message: Option<KafkaString>,
21    /// The response from the describe quorum API.
22    pub topics: Vec<TopicData>,
23    /// The nodes in the quorum.
24    pub nodes: Vec<Node>,
25    pub _unknown_tagged_fields: Vec<RawTaggedField>,
26}
27impl Default for DescribeQuorumResponseData {
28    fn default() -> Self {
29        Self {
30            error_code: 0_i16,
31            error_message: None,
32            topics: Vec::new(),
33            nodes: Vec::new(),
34            _unknown_tagged_fields: Vec::new(),
35        }
36    }
37}
38impl DescribeQuorumResponseData {
39    pub fn with_error_code(mut self, value: i16) -> Self {
40        self.error_code = value;
41        self
42    }
43    pub fn with_error_message(mut self, value: Option<KafkaString>) -> Self {
44        self.error_message = value;
45        self
46    }
47    pub fn with_topics(mut self, value: Vec<TopicData>) -> Self {
48        self.topics = value;
49        self
50    }
51    pub fn with_nodes(mut self, value: Vec<Node>) -> Self {
52        self.nodes = value;
53        self
54    }
55    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
56        if version < 0 || version > 2 {
57            return Err(UnsupportedVersion::new(55, version).into());
58        }
59        let error_code;
60        let mut error_message = None;
61        let topics;
62        let mut nodes = Vec::new();
63        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
64        error_code = read_i16(buf)?;
65        if version >= 2 {
66            error_message = read_compact_nullable_string(buf)?;
67        }
68        topics = {
69            let len = read_compact_array_length(buf)?;
70            let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
71            for _ in 0..len {
72                arr.push(TopicData::read(buf, version)?);
73            }
74            arr
75        };
76        if version >= 2 {
77            nodes = {
78                let len = read_compact_array_length(buf)?;
79                let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
80                for _ in 0..len {
81                    arr.push(Node::read(buf, version)?);
82                }
83                arr
84            };
85        }
86        let tagged_fields = read_tagged_fields(buf)?;
87        for field in &tagged_fields {
88            match field.tag {
89                _ => {
90                    _unknown_tagged_fields.push(field.clone());
91                },
92            }
93        }
94        Ok(Self {
95            error_code,
96            error_message,
97            topics,
98            nodes,
99            _unknown_tagged_fields,
100        })
101    }
102    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
103        if version < 0 || version > 2 {
104            return Err(UnsupportedVersion::new(55, version).into());
105        }
106        write_i16(buf, self.error_code);
107        if version >= 2 {
108            write_compact_nullable_string(buf, self.error_message.as_ref())?;
109        } else if self.error_message != None {
110            return Err(UnsupportedFieldVersion::new(55, "error_message", version).into());
111        }
112        write_compact_array_length(buf, self.topics.len() as i32);
113        for el in &self.topics {
114            el.write(buf, version)?;
115        }
116        if version >= 2 {
117            write_compact_array_length(buf, self.nodes.len() as i32);
118            for el in &self.nodes {
119                el.write(buf, version)?;
120            }
121        } else if self.nodes != Vec::new() {
122            return Err(UnsupportedFieldVersion::new(55, "nodes", version).into());
123        }
124        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
125        all_tags.sort_by_key(|f| f.tag);
126        write_tagged_fields(buf, &all_tags)?;
127        Ok(())
128    }
129    pub fn encoded_len(&self, version: i16) -> Result<usize> {
130        if version < 0 || version > 2 {
131            return Err(UnsupportedVersion::new(55, version).into());
132        }
133        let mut len: usize = 0;
134        len += 2;
135        if version >= 2 {
136            len += compact_nullable_string_len(self.error_message.as_ref())?;
137        } else if self.error_message != None {
138            return Err(UnsupportedFieldVersion::new(55, "error_message", version).into());
139        }
140        len += compact_array_length_len(self.topics.len() as i32);
141        for el in &self.topics {
142            len += el.encoded_len(version)?;
143        }
144        if version >= 2 {
145            len += compact_array_length_len(self.nodes.len() as i32);
146            for el in &self.nodes {
147                len += el.encoded_len(version)?;
148            }
149        } else if self.nodes != Vec::new() {
150            return Err(UnsupportedFieldVersion::new(55, "nodes", version).into());
151        }
152        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
153        all_tags.sort_by_key(|f| f.tag);
154        len += tagged_fields_len(&all_tags)?;
155        Ok(len)
156    }
157}
158#[derive(Debug, Clone, PartialEq)]
159pub struct TopicData {
160    /// The topic name.
161    pub topic_name: KafkaString,
162    /// The partition data.
163    pub partitions: Vec<PartitionData>,
164    pub _unknown_tagged_fields: Vec<RawTaggedField>,
165}
166impl Default for TopicData {
167    fn default() -> Self {
168        Self {
169            topic_name: KafkaString::default(),
170            partitions: Vec::new(),
171            _unknown_tagged_fields: Vec::new(),
172        }
173    }
174}
175impl TopicData {
176    pub fn with_topic_name(mut self, value: KafkaString) -> Self {
177        self.topic_name = value;
178        self
179    }
180    pub fn with_partitions(mut self, value: Vec<PartitionData>) -> Self {
181        self.partitions = value;
182        self
183    }
184    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
185        let topic_name;
186        let partitions;
187        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
188        topic_name = read_compact_string(buf)?;
189        partitions = {
190            let len = read_compact_array_length(buf)?;
191            let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
192            for _ in 0..len {
193                arr.push(PartitionData::read(buf, version)?);
194            }
195            arr
196        };
197        let tagged_fields = read_tagged_fields(buf)?;
198        for field in &tagged_fields {
199            match field.tag {
200                _ => {
201                    _unknown_tagged_fields.push(field.clone());
202                },
203            }
204        }
205        Ok(Self {
206            topic_name,
207            partitions,
208            _unknown_tagged_fields,
209        })
210    }
211    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
212        write_compact_string(buf, &self.topic_name)?;
213        write_compact_array_length(buf, self.partitions.len() as i32);
214        for el in &self.partitions {
215            el.write(buf, version)?;
216        }
217        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
218        all_tags.sort_by_key(|f| f.tag);
219        write_tagged_fields(buf, &all_tags)?;
220        Ok(())
221    }
222    pub fn encoded_len(&self, version: i16) -> Result<usize> {
223        let mut len: usize = 0;
224        len += compact_string_len(&self.topic_name)?;
225        len += compact_array_length_len(self.partitions.len() as i32);
226        for el in &self.partitions {
227            len += el.encoded_len(version)?;
228        }
229        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
230        all_tags.sort_by_key(|f| f.tag);
231        len += tagged_fields_len(&all_tags)?;
232        Ok(len)
233    }
234}
235#[derive(Debug, Clone, PartialEq)]
236pub struct PartitionData {
237    /// The partition index.
238    pub partition_index: i32,
239    /// The partition error code.
240    pub error_code: i16,
241    /// The error message, or null if there was no error.
242    pub error_message: Option<KafkaString>,
243    /// The ID of the current leader or -1 if the leader is unknown.
244    pub leader_id: i32,
245    /// The latest known leader epoch.
246    pub leader_epoch: i32,
247    /// The high water mark.
248    pub high_watermark: i64,
249    /// The current voters of the partition.
250    pub current_voters: Vec<ReplicaState>,
251    /// The observers of the partition.
252    pub observers: Vec<ReplicaState>,
253    pub _unknown_tagged_fields: Vec<RawTaggedField>,
254}
255impl Default for PartitionData {
256    fn default() -> Self {
257        Self {
258            partition_index: 0_i32,
259            error_code: 0_i16,
260            error_message: None,
261            leader_id: 0_i32,
262            leader_epoch: 0_i32,
263            high_watermark: 0_i64,
264            current_voters: Vec::new(),
265            observers: Vec::new(),
266            _unknown_tagged_fields: Vec::new(),
267        }
268    }
269}
270impl PartitionData {
271    pub fn with_partition_index(mut self, value: i32) -> Self {
272        self.partition_index = value;
273        self
274    }
275    pub fn with_error_code(mut self, value: i16) -> Self {
276        self.error_code = value;
277        self
278    }
279    pub fn with_error_message(mut self, value: Option<KafkaString>) -> Self {
280        self.error_message = value;
281        self
282    }
283    pub fn with_leader_id(mut self, value: i32) -> Self {
284        self.leader_id = value;
285        self
286    }
287    pub fn with_leader_epoch(mut self, value: i32) -> Self {
288        self.leader_epoch = value;
289        self
290    }
291    pub fn with_high_watermark(mut self, value: i64) -> Self {
292        self.high_watermark = value;
293        self
294    }
295    pub fn with_current_voters(mut self, value: Vec<ReplicaState>) -> Self {
296        self.current_voters = value;
297        self
298    }
299    pub fn with_observers(mut self, value: Vec<ReplicaState>) -> Self {
300        self.observers = value;
301        self
302    }
303    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
304        let partition_index;
305        let error_code;
306        let mut error_message = None;
307        let leader_id;
308        let leader_epoch;
309        let high_watermark;
310        let current_voters;
311        let observers;
312        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
313        partition_index = read_i32(buf)?;
314        error_code = read_i16(buf)?;
315        if version >= 2 {
316            error_message = read_compact_nullable_string(buf)?;
317        }
318        leader_id = read_i32(buf)?;
319        leader_epoch = read_i32(buf)?;
320        high_watermark = read_i64(buf)?;
321        current_voters = {
322            let len = read_compact_array_length(buf)?;
323            let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
324            for _ in 0..len {
325                arr.push(ReplicaState::read(buf, version)?);
326            }
327            arr
328        };
329        observers = {
330            let len = read_compact_array_length(buf)?;
331            let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
332            for _ in 0..len {
333                arr.push(ReplicaState::read(buf, version)?);
334            }
335            arr
336        };
337        let tagged_fields = read_tagged_fields(buf)?;
338        for field in &tagged_fields {
339            match field.tag {
340                _ => {
341                    _unknown_tagged_fields.push(field.clone());
342                },
343            }
344        }
345        Ok(Self {
346            partition_index,
347            error_code,
348            error_message,
349            leader_id,
350            leader_epoch,
351            high_watermark,
352            current_voters,
353            observers,
354            _unknown_tagged_fields,
355        })
356    }
357    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
358        write_i32(buf, self.partition_index);
359        write_i16(buf, self.error_code);
360        if version >= 2 {
361            write_compact_nullable_string(buf, self.error_message.as_ref())?;
362        } else if self.error_message != None {
363            return Err(UnsupportedFieldVersion::new(55, "error_message", version).into());
364        }
365        write_i32(buf, self.leader_id);
366        write_i32(buf, self.leader_epoch);
367        write_i64(buf, self.high_watermark);
368        write_compact_array_length(buf, self.current_voters.len() as i32);
369        for el in &self.current_voters {
370            el.write(buf, version)?;
371        }
372        write_compact_array_length(buf, self.observers.len() as i32);
373        for el in &self.observers {
374            el.write(buf, version)?;
375        }
376        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
377        all_tags.sort_by_key(|f| f.tag);
378        write_tagged_fields(buf, &all_tags)?;
379        Ok(())
380    }
381    pub fn encoded_len(&self, version: i16) -> Result<usize> {
382        let mut len: usize = 0;
383        len += 4;
384        len += 2;
385        if version >= 2 {
386            len += compact_nullable_string_len(self.error_message.as_ref())?;
387        } else if self.error_message != None {
388            return Err(UnsupportedFieldVersion::new(55, "error_message", version).into());
389        }
390        len += 4;
391        len += 4;
392        len += 8;
393        len += compact_array_length_len(self.current_voters.len() as i32);
394        for el in &self.current_voters {
395            len += el.encoded_len(version)?;
396        }
397        len += compact_array_length_len(self.observers.len() as i32);
398        for el in &self.observers {
399            len += el.encoded_len(version)?;
400        }
401        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
402        all_tags.sort_by_key(|f| f.tag);
403        len += tagged_fields_len(&all_tags)?;
404        Ok(len)
405    }
406}
407#[derive(Debug, Clone, PartialEq)]
408pub struct Node {
409    /// The ID of the associated node.
410    pub node_id: i32,
411    /// The listeners of this controller.
412    pub listeners: Vec<Listener>,
413    pub _unknown_tagged_fields: Vec<RawTaggedField>,
414}
415impl Default for Node {
416    fn default() -> Self {
417        Self {
418            node_id: 0_i32,
419            listeners: Vec::new(),
420            _unknown_tagged_fields: Vec::new(),
421        }
422    }
423}
424impl Node {
425    pub fn with_node_id(mut self, value: i32) -> Self {
426        self.node_id = value;
427        self
428    }
429    pub fn with_listeners(mut self, value: Vec<Listener>) -> Self {
430        self.listeners = value;
431        self
432    }
433    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
434        let node_id;
435        let listeners;
436        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
437        node_id = read_i32(buf)?;
438        listeners = {
439            let len = read_compact_array_length(buf)?;
440            let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
441            for _ in 0..len {
442                arr.push(Listener::read(buf, version)?);
443            }
444            arr
445        };
446        let tagged_fields = read_tagged_fields(buf)?;
447        for field in &tagged_fields {
448            match field.tag {
449                _ => {
450                    _unknown_tagged_fields.push(field.clone());
451                },
452            }
453        }
454        Ok(Self {
455            node_id,
456            listeners,
457            _unknown_tagged_fields,
458        })
459    }
460    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
461        write_i32(buf, self.node_id);
462        write_compact_array_length(buf, self.listeners.len() as i32);
463        for el in &self.listeners {
464            el.write(buf, version)?;
465        }
466        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
467        all_tags.sort_by_key(|f| f.tag);
468        write_tagged_fields(buf, &all_tags)?;
469        Ok(())
470    }
471    pub fn encoded_len(&self, version: i16) -> Result<usize> {
472        let mut len: usize = 0;
473        len += 4;
474        len += compact_array_length_len(self.listeners.len() as i32);
475        for el in &self.listeners {
476            len += el.encoded_len(version)?;
477        }
478        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
479        all_tags.sort_by_key(|f| f.tag);
480        len += tagged_fields_len(&all_tags)?;
481        Ok(len)
482    }
483}
484#[derive(Debug, Clone, PartialEq)]
485pub struct Listener {
486    /// The name of the endpoint.
487    pub name: KafkaString,
488    /// The hostname.
489    pub host: KafkaString,
490    /// The port.
491    pub port: u16,
492    pub _unknown_tagged_fields: Vec<RawTaggedField>,
493}
494impl Default for Listener {
495    fn default() -> Self {
496        Self {
497            name: KafkaString::default(),
498            host: KafkaString::default(),
499            port: 0_u16,
500            _unknown_tagged_fields: Vec::new(),
501        }
502    }
503}
504impl Listener {
505    pub fn with_name(mut self, value: KafkaString) -> Self {
506        self.name = value;
507        self
508    }
509    pub fn with_host(mut self, value: KafkaString) -> Self {
510        self.host = value;
511        self
512    }
513    pub fn with_port(mut self, value: u16) -> Self {
514        self.port = value;
515        self
516    }
517    pub fn read(buf: &mut Bytes, _version: i16) -> Result<Self> {
518        let name;
519        let host;
520        let port;
521        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
522        name = read_compact_string(buf)?;
523        host = read_compact_string(buf)?;
524        port = read_u16(buf)?;
525        let tagged_fields = read_tagged_fields(buf)?;
526        for field in &tagged_fields {
527            match field.tag {
528                _ => {
529                    _unknown_tagged_fields.push(field.clone());
530                },
531            }
532        }
533        Ok(Self {
534            name,
535            host,
536            port,
537            _unknown_tagged_fields,
538        })
539    }
540    pub fn write(&self, buf: &mut BytesMut, _version: i16) -> Result<()> {
541        write_compact_string(buf, &self.name)?;
542        write_compact_string(buf, &self.host)?;
543        write_u16(buf, self.port);
544        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
545        all_tags.sort_by_key(|f| f.tag);
546        write_tagged_fields(buf, &all_tags)?;
547        Ok(())
548    }
549    pub fn encoded_len(&self, _version: i16) -> Result<usize> {
550        let mut len: usize = 0;
551        len += compact_string_len(&self.name)?;
552        len += compact_string_len(&self.host)?;
553        len += 2;
554        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
555        all_tags.sort_by_key(|f| f.tag);
556        len += tagged_fields_len(&all_tags)?;
557        Ok(len)
558    }
559}
560#[derive(Debug, Clone, PartialEq)]
561pub struct ReplicaState {
562    /// The ID of the replica.
563    pub replica_id: i32,
564    /// The replica directory ID of the replica.
565    pub replica_directory_id: KafkaUuid,
566    /// The last known log end offset of the follower or -1 if it is unknown.
567    pub log_end_offset: i64,
568    /// The last known leader wall clock time time when a follower fetched from the leader. This is
569    /// reported as -1 both for the current leader or if it is unknown for a voter.
570    pub last_fetch_timestamp: i64,
571    /// The leader wall clock append time of the offset for which the follower made the most recent
572    /// fetch request. This is reported as the current time for the leader and -1 if unknown for a
573    /// voter.
574    pub last_caught_up_timestamp: i64,
575    pub _unknown_tagged_fields: Vec<RawTaggedField>,
576}
577impl Default for ReplicaState {
578    fn default() -> Self {
579        Self {
580            replica_id: 0_i32,
581            replica_directory_id: KafkaUuid::ZERO,
582            log_end_offset: 0_i64,
583            last_fetch_timestamp: -1i64,
584            last_caught_up_timestamp: -1i64,
585            _unknown_tagged_fields: Vec::new(),
586        }
587    }
588}
589impl ReplicaState {
590    pub fn with_replica_id(mut self, value: i32) -> Self {
591        self.replica_id = value;
592        self
593    }
594    pub fn with_replica_directory_id(mut self, value: KafkaUuid) -> Self {
595        self.replica_directory_id = value;
596        self
597    }
598    pub fn with_log_end_offset(mut self, value: i64) -> Self {
599        self.log_end_offset = value;
600        self
601    }
602    pub fn with_last_fetch_timestamp(mut self, value: i64) -> Self {
603        self.last_fetch_timestamp = value;
604        self
605    }
606    pub fn with_last_caught_up_timestamp(mut self, value: i64) -> Self {
607        self.last_caught_up_timestamp = value;
608        self
609    }
610    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
611        let replica_id;
612        let mut replica_directory_id = KafkaUuid::ZERO;
613        let log_end_offset;
614        let mut last_fetch_timestamp = -1i64;
615        let mut last_caught_up_timestamp = -1i64;
616        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
617        replica_id = read_i32(buf)?;
618        if version >= 2 {
619            replica_directory_id = read_uuid(buf)?;
620        }
621        log_end_offset = read_i64(buf)?;
622        if version >= 1 {
623            last_fetch_timestamp = read_i64(buf)?;
624        }
625        if version >= 1 {
626            last_caught_up_timestamp = read_i64(buf)?;
627        }
628        let tagged_fields = read_tagged_fields(buf)?;
629        for field in &tagged_fields {
630            match field.tag {
631                _ => {
632                    _unknown_tagged_fields.push(field.clone());
633                },
634            }
635        }
636        Ok(Self {
637            replica_id,
638            replica_directory_id,
639            log_end_offset,
640            last_fetch_timestamp,
641            last_caught_up_timestamp,
642            _unknown_tagged_fields,
643        })
644    }
645    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
646        write_i32(buf, self.replica_id);
647        if version >= 2 {
648            write_uuid(buf, &self.replica_directory_id);
649        } else if self.replica_directory_id != KafkaUuid::ZERO {
650            return Err(UnsupportedFieldVersion::new(55, "replica_directory_id", version).into());
651        }
652        write_i64(buf, self.log_end_offset);
653        if version >= 1 {
654            write_i64(buf, self.last_fetch_timestamp);
655        } else if self.last_fetch_timestamp != -1i64 {
656            return Err(UnsupportedFieldVersion::new(55, "last_fetch_timestamp", version).into());
657        }
658        if version >= 1 {
659            write_i64(buf, self.last_caught_up_timestamp);
660        } else if self.last_caught_up_timestamp != -1i64 {
661            return Err(
662                UnsupportedFieldVersion::new(55, "last_caught_up_timestamp", version).into(),
663            );
664        }
665        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
666        all_tags.sort_by_key(|f| f.tag);
667        write_tagged_fields(buf, &all_tags)?;
668        Ok(())
669    }
670    pub fn encoded_len(&self, version: i16) -> Result<usize> {
671        let mut len: usize = 0;
672        len += 4;
673        if version >= 2 {
674            len += 16;
675        } else if self.replica_directory_id != KafkaUuid::ZERO {
676            return Err(UnsupportedFieldVersion::new(55, "replica_directory_id", version).into());
677        }
678        len += 8;
679        if version >= 1 {
680            len += 8;
681        } else if self.last_fetch_timestamp != -1i64 {
682            return Err(UnsupportedFieldVersion::new(55, "last_fetch_timestamp", version).into());
683        }
684        if version >= 1 {
685            len += 8;
686        } else if self.last_caught_up_timestamp != -1i64 {
687            return Err(
688                UnsupportedFieldVersion::new(55, "last_caught_up_timestamp", version).into(),
689            );
690        }
691        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
692        all_tags.sort_by_key(|f| f.tag);
693        len += tagged_fields_len(&all_tags)?;
694        Ok(len)
695    }
696}