Skip to main content

kacrab_protocol/generated/
offset_fetch_request.rs

1//! Generated from OffsetFetchRequest.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 OffsetFetchRequestData {
17    /// The group to fetch offsets for.
18    pub group_id: KafkaString,
19    /// Each topic we would like to fetch offsets for, or null to fetch offsets for all topics.
20    pub topics: Option<Vec<OffsetFetchRequestTopic>>,
21    /// Each group we would like to fetch offsets for.
22    pub groups: Vec<OffsetFetchRequestGroup>,
23    /// Whether broker should hold on returning unstable offsets but set a retriable error code for
24    /// the partitions.
25    pub require_stable: bool,
26    pub _unknown_tagged_fields: Vec<RawTaggedField>,
27}
28impl Default for OffsetFetchRequestData {
29    fn default() -> Self {
30        Self {
31            group_id: KafkaString::default(),
32            topics: None,
33            groups: Vec::new(),
34            require_stable: false,
35            _unknown_tagged_fields: Vec::new(),
36        }
37    }
38}
39impl OffsetFetchRequestData {
40    pub fn with_group_id(mut self, value: KafkaString) -> Self {
41        self.group_id = value;
42        self
43    }
44    pub fn with_topics(mut self, value: Option<Vec<OffsetFetchRequestTopic>>) -> Self {
45        self.topics = value;
46        self
47    }
48    pub fn with_groups(mut self, value: Vec<OffsetFetchRequestGroup>) -> Self {
49        self.groups = value;
50        self
51    }
52    pub fn with_require_stable(mut self, value: bool) -> Self {
53        self.require_stable = value;
54        self
55    }
56    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
57        if version < 1 || version > 10 {
58            return Err(UnsupportedVersion::new(9, version).into());
59        }
60        let mut group_id = KafkaString::default();
61        let mut topics = None;
62        let mut groups = Vec::new();
63        let mut require_stable = false;
64        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
65        if version <= 7 {
66            if version >= 6 {
67                group_id = read_compact_string(buf)?;
68            } else {
69                group_id = read_string(buf)?;
70            }
71        }
72        if version <= 7 {
73            if version >= 2 {
74                if version >= 6 {
75                    topics = {
76                        let len = read_compact_array_length(buf)?;
77                        if len < 0 {
78                            None
79                        } else {
80                            let mut arr = Vec::with_capacity(len as usize);
81                            for _ in 0..len {
82                                arr.push(OffsetFetchRequestTopic::read(buf, version)?);
83                            }
84                            Some(arr)
85                        }
86                    };
87                } else {
88                    topics = {
89                        let len = read_array_length(buf)?;
90                        if len < 0 {
91                            None
92                        } else {
93                            let mut arr = Vec::with_capacity(len as usize);
94                            for _ in 0..len {
95                                arr.push(OffsetFetchRequestTopic::read(buf, version)?);
96                            }
97                            Some(arr)
98                        }
99                    };
100                }
101            } else {
102                topics = Some({
103                    let len = read_array_length(buf)?;
104                    let mut arr = Vec::with_capacity(len.max(0) as usize);
105                    for _ in 0..len {
106                        arr.push(OffsetFetchRequestTopic::read(buf, version)?);
107                    }
108                    arr
109                });
110            }
111        }
112        if version >= 8 {
113            groups = {
114                let len = read_compact_array_length(buf)?;
115                let mut arr = Vec::with_capacity(len.max(0) as usize);
116                for _ in 0..len {
117                    arr.push(OffsetFetchRequestGroup::read(buf, version)?);
118                }
119                arr
120            };
121        }
122        if version >= 7 {
123            require_stable = read_bool(buf)?;
124        }
125        if version >= 6 {
126            let tagged_fields = read_tagged_fields(buf)?;
127            for field in &tagged_fields {
128                match field.tag {
129                    _ => {
130                        _unknown_tagged_fields.push(field.clone());
131                    },
132                }
133            }
134        }
135        Ok(Self {
136            group_id,
137            topics,
138            groups,
139            require_stable,
140            _unknown_tagged_fields,
141        })
142    }
143    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
144        if version < 1 || version > 10 {
145            return Err(UnsupportedVersion::new(9, version).into());
146        }
147        if version <= 7 {
148            if version >= 6 {
149                write_compact_string(buf, &self.group_id)?;
150            } else {
151                write_string(buf, &self.group_id)?;
152            }
153        } else if self.group_id != KafkaString::default() {
154            return Err(UnsupportedFieldVersion::new(9, "group_id", version).into());
155        }
156        if version <= 7 {
157            if version >= 2 {
158                if version >= 6 {
159                    match &self.topics {
160                        None => {
161                            write_compact_array_length(buf, -1);
162                        },
163                        Some(arr) => {
164                            write_compact_array_length(buf, arr.len() as i32);
165                            for el in arr {
166                                el.write(buf, version)?;
167                            }
168                        },
169                    }
170                } else {
171                    match &self.topics {
172                        None => {
173                            write_array_length(buf, -1);
174                        },
175                        Some(arr) => {
176                            write_array_length(buf, arr.len() as i32);
177                            for el in arr {
178                                el.write(buf, version)?;
179                            }
180                        },
181                    }
182                }
183            } else {
184                match &self.topics {
185                    Some(arr) => {
186                        write_array_length(buf, arr.len() as i32);
187                        for el in arr {
188                            el.write(buf, version)?;
189                        }
190                    },
191                    None => {
192                        write_array_length(buf, 0);
193                    },
194                }
195            }
196        } else if self.topics != None {
197            return Err(UnsupportedFieldVersion::new(9, "topics", version).into());
198        }
199        if version >= 8 {
200            write_compact_array_length(buf, self.groups.len() as i32);
201            for el in &self.groups {
202                el.write(buf, version)?;
203            }
204        } else if self.groups != Vec::new() {
205            return Err(UnsupportedFieldVersion::new(9, "groups", version).into());
206        }
207        if version >= 7 {
208            write_bool(buf, self.require_stable);
209        } else if self.require_stable != false {
210            return Err(UnsupportedFieldVersion::new(9, "require_stable", version).into());
211        }
212        if version >= 6 {
213            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
214            all_tags.sort_by_key(|f| f.tag);
215            write_tagged_fields(buf, &all_tags)?;
216        }
217        Ok(())
218    }
219    pub fn encoded_len(&self, version: i16) -> Result<usize> {
220        if version < 1 || version > 10 {
221            return Err(UnsupportedVersion::new(9, version).into());
222        }
223        let mut len: usize = 0;
224        if version <= 7 {
225            if version >= 6 {
226                len += compact_string_len(&self.group_id)?;
227            } else {
228                len += string_len(&self.group_id)?;
229            }
230        } else if self.group_id != KafkaString::default() {
231            return Err(UnsupportedFieldVersion::new(9, "group_id", version).into());
232        }
233        if version <= 7 {
234            if version >= 2 {
235                if version >= 6 {
236                    match &self.topics {
237                        None => {
238                            len += compact_array_length_len(-1);
239                        },
240                        Some(arr) => {
241                            len += compact_array_length_len(arr.len() as i32);
242                            for el in arr {
243                                len += el.encoded_len(version)?;
244                            }
245                        },
246                    }
247                } else {
248                    match &self.topics {
249                        None => {
250                            len += array_length_len();
251                        },
252                        Some(arr) => {
253                            len += array_length_len();
254                            for el in arr {
255                                len += el.encoded_len(version)?;
256                            }
257                        },
258                    }
259                }
260            } else {
261                match &self.topics {
262                    Some(arr) => {
263                        len += array_length_len();
264                        for el in arr {
265                            len += el.encoded_len(version)?;
266                        }
267                    },
268                    None => {
269                        len += array_length_len();
270                    },
271                }
272            }
273        } else if self.topics != None {
274            return Err(UnsupportedFieldVersion::new(9, "topics", version).into());
275        }
276        if version >= 8 {
277            len += compact_array_length_len(self.groups.len() as i32);
278            for el in &self.groups {
279                len += el.encoded_len(version)?;
280            }
281        } else if self.groups != Vec::new() {
282            return Err(UnsupportedFieldVersion::new(9, "groups", version).into());
283        }
284        if version >= 7 {
285            len += 1;
286        } else if self.require_stable != false {
287            return Err(UnsupportedFieldVersion::new(9, "require_stable", version).into());
288        }
289        if version >= 6 {
290            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
291            all_tags.sort_by_key(|f| f.tag);
292            len += tagged_fields_len(&all_tags)?;
293        }
294        Ok(len)
295    }
296}
297#[derive(Debug, Clone, PartialEq)]
298pub struct OffsetFetchRequestTopic {
299    /// The topic name.
300    pub name: KafkaString,
301    /// The partition indexes we would like to fetch offsets for.
302    pub partition_indexes: Vec<i32>,
303    pub _unknown_tagged_fields: Vec<RawTaggedField>,
304}
305impl Default for OffsetFetchRequestTopic {
306    fn default() -> Self {
307        Self {
308            name: KafkaString::default(),
309            partition_indexes: Vec::new(),
310            _unknown_tagged_fields: Vec::new(),
311        }
312    }
313}
314impl OffsetFetchRequestTopic {
315    pub fn with_name(mut self, value: KafkaString) -> Self {
316        self.name = value;
317        self
318    }
319    pub fn with_partition_indexes(mut self, value: Vec<i32>) -> Self {
320        self.partition_indexes = value;
321        self
322    }
323    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
324        let name;
325        let partition_indexes;
326        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
327        if version >= 6 {
328            name = read_compact_string(buf)?;
329        } else {
330            name = read_string(buf)?;
331        }
332        if version >= 6 {
333            partition_indexes = {
334                let len = read_compact_array_length(buf)?;
335                let mut arr = Vec::with_capacity(len.max(0) as usize);
336                for _ in 0..len {
337                    arr.push(read_i32(buf)?);
338                }
339                arr
340            };
341        } else {
342            partition_indexes = {
343                let len = read_array_length(buf)?;
344                let mut arr = Vec::with_capacity(len.max(0) as usize);
345                for _ in 0..len {
346                    arr.push(read_i32(buf)?);
347                }
348                arr
349            };
350        }
351        if version >= 6 {
352            let tagged_fields = read_tagged_fields(buf)?;
353            for field in &tagged_fields {
354                match field.tag {
355                    _ => {
356                        _unknown_tagged_fields.push(field.clone());
357                    },
358                }
359            }
360        }
361        Ok(Self {
362            name,
363            partition_indexes,
364            _unknown_tagged_fields,
365        })
366    }
367    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
368        if version >= 6 {
369            write_compact_string(buf, &self.name)?;
370        } else {
371            write_string(buf, &self.name)?;
372        }
373        if version >= 6 {
374            write_compact_array_length(buf, self.partition_indexes.len() as i32);
375            for el in &self.partition_indexes {
376                write_i32(buf, *el);
377            }
378        } else {
379            write_array_length(buf, self.partition_indexes.len() as i32);
380            for el in &self.partition_indexes {
381                write_i32(buf, *el);
382            }
383        }
384        if version >= 6 {
385            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
386            all_tags.sort_by_key(|f| f.tag);
387            write_tagged_fields(buf, &all_tags)?;
388        }
389        Ok(())
390    }
391    pub fn encoded_len(&self, version: i16) -> Result<usize> {
392        let mut len: usize = 0;
393        if version >= 6 {
394            len += compact_string_len(&self.name)?;
395        } else {
396            len += string_len(&self.name)?;
397        }
398        if version >= 6 {
399            len += compact_array_length_len(self.partition_indexes.len() as i32);
400            len += self.partition_indexes.len() * 4usize;
401        } else {
402            len += array_length_len();
403            len += self.partition_indexes.len() * 4usize;
404        }
405        if version >= 6 {
406            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
407            all_tags.sort_by_key(|f| f.tag);
408            len += tagged_fields_len(&all_tags)?;
409        }
410        Ok(len)
411    }
412}
413#[derive(Debug, Clone, PartialEq)]
414pub struct OffsetFetchRequestGroup {
415    /// The group ID.
416    pub group_id: KafkaString,
417    /// The member id.
418    pub member_id: Option<KafkaString>,
419    /// The member epoch if using the new consumer protocol (KIP-848).
420    pub member_epoch: i32,
421    /// Each topic we would like to fetch offsets for, or null to fetch offsets for all topics.
422    pub topics: Option<Vec<OffsetFetchRequestTopics>>,
423    pub _unknown_tagged_fields: Vec<RawTaggedField>,
424}
425impl Default for OffsetFetchRequestGroup {
426    fn default() -> Self {
427        Self {
428            group_id: KafkaString::default(),
429            member_id: None,
430            member_epoch: -1i32,
431            topics: None,
432            _unknown_tagged_fields: Vec::new(),
433        }
434    }
435}
436impl OffsetFetchRequestGroup {
437    pub fn with_group_id(mut self, value: KafkaString) -> Self {
438        self.group_id = value;
439        self
440    }
441    pub fn with_member_id(mut self, value: Option<KafkaString>) -> Self {
442        self.member_id = value;
443        self
444    }
445    pub fn with_member_epoch(mut self, value: i32) -> Self {
446        self.member_epoch = value;
447        self
448    }
449    pub fn with_topics(mut self, value: Option<Vec<OffsetFetchRequestTopics>>) -> Self {
450        self.topics = value;
451        self
452    }
453    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
454        let group_id;
455        let mut member_id = None;
456        let mut member_epoch = -1i32;
457        let topics;
458        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
459        group_id = read_compact_string(buf)?;
460        if version >= 9 {
461            member_id = read_compact_nullable_string(buf)?;
462        }
463        if version >= 9 {
464            member_epoch = read_i32(buf)?;
465        }
466        topics = {
467            let len = read_compact_array_length(buf)?;
468            if len < 0 {
469                None
470            } else {
471                let mut arr = Vec::with_capacity(len as usize);
472                for _ in 0..len {
473                    arr.push(OffsetFetchRequestTopics::read(buf, version)?);
474                }
475                Some(arr)
476            }
477        };
478        let tagged_fields = read_tagged_fields(buf)?;
479        for field in &tagged_fields {
480            match field.tag {
481                _ => {
482                    _unknown_tagged_fields.push(field.clone());
483                },
484            }
485        }
486        Ok(Self {
487            group_id,
488            member_id,
489            member_epoch,
490            topics,
491            _unknown_tagged_fields,
492        })
493    }
494    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
495        write_compact_string(buf, &self.group_id)?;
496        if version >= 9 {
497            write_compact_nullable_string(buf, self.member_id.as_ref())?;
498        } else if self.member_id != None {
499            return Err(UnsupportedFieldVersion::new(9, "member_id", version).into());
500        }
501        if version >= 9 {
502            write_i32(buf, self.member_epoch);
503        } else if self.member_epoch != -1i32 {
504            return Err(UnsupportedFieldVersion::new(9, "member_epoch", version).into());
505        }
506        match &self.topics {
507            None => {
508                write_compact_array_length(buf, -1);
509            },
510            Some(arr) => {
511                write_compact_array_length(buf, arr.len() as i32);
512                for el in arr {
513                    el.write(buf, version)?;
514                }
515            },
516        }
517        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
518        all_tags.sort_by_key(|f| f.tag);
519        write_tagged_fields(buf, &all_tags)?;
520        Ok(())
521    }
522    pub fn encoded_len(&self, version: i16) -> Result<usize> {
523        let mut len: usize = 0;
524        len += compact_string_len(&self.group_id)?;
525        if version >= 9 {
526            len += compact_nullable_string_len(self.member_id.as_ref())?;
527        } else if self.member_id != None {
528            return Err(UnsupportedFieldVersion::new(9, "member_id", version).into());
529        }
530        if version >= 9 {
531            len += 4;
532        } else if self.member_epoch != -1i32 {
533            return Err(UnsupportedFieldVersion::new(9, "member_epoch", version).into());
534        }
535        match &self.topics {
536            None => {
537                len += compact_array_length_len(-1);
538            },
539            Some(arr) => {
540                len += compact_array_length_len(arr.len() as i32);
541                for el in arr {
542                    len += el.encoded_len(version)?;
543                }
544            },
545        }
546        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
547        all_tags.sort_by_key(|f| f.tag);
548        len += tagged_fields_len(&all_tags)?;
549        Ok(len)
550    }
551}
552#[derive(Debug, Clone, PartialEq)]
553pub struct OffsetFetchRequestTopics {
554    /// The topic name.
555    pub name: KafkaString,
556    /// The topic ID.
557    pub topic_id: KafkaUuid,
558    /// The partition indexes we would like to fetch offsets for.
559    pub partition_indexes: Vec<i32>,
560    pub _unknown_tagged_fields: Vec<RawTaggedField>,
561}
562impl Default for OffsetFetchRequestTopics {
563    fn default() -> Self {
564        Self {
565            name: KafkaString::default(),
566            topic_id: KafkaUuid::ZERO,
567            partition_indexes: Vec::new(),
568            _unknown_tagged_fields: Vec::new(),
569        }
570    }
571}
572impl OffsetFetchRequestTopics {
573    pub fn with_name(mut self, value: KafkaString) -> Self {
574        self.name = value;
575        self
576    }
577    pub fn with_topic_id(mut self, value: KafkaUuid) -> Self {
578        self.topic_id = value;
579        self
580    }
581    pub fn with_partition_indexes(mut self, value: Vec<i32>) -> Self {
582        self.partition_indexes = value;
583        self
584    }
585    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
586        let mut name = KafkaString::default();
587        let mut topic_id = KafkaUuid::ZERO;
588        let partition_indexes;
589        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
590        if version <= 9 {
591            name = read_compact_string(buf)?;
592        }
593        if version >= 10 {
594            topic_id = read_uuid(buf)?;
595        }
596        partition_indexes = {
597            let len = read_compact_array_length(buf)?;
598            let mut arr = Vec::with_capacity(len.max(0) as usize);
599            for _ in 0..len {
600                arr.push(read_i32(buf)?);
601            }
602            arr
603        };
604        let tagged_fields = read_tagged_fields(buf)?;
605        for field in &tagged_fields {
606            match field.tag {
607                _ => {
608                    _unknown_tagged_fields.push(field.clone());
609                },
610            }
611        }
612        Ok(Self {
613            name,
614            topic_id,
615            partition_indexes,
616            _unknown_tagged_fields,
617        })
618    }
619    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
620        if version <= 9 {
621            write_compact_string(buf, &self.name)?;
622        } else if self.name != KafkaString::default() {
623            return Err(UnsupportedFieldVersion::new(9, "name", version).into());
624        }
625        if version >= 10 {
626            write_uuid(buf, &self.topic_id);
627        } else if self.topic_id != KafkaUuid::ZERO {
628            return Err(UnsupportedFieldVersion::new(9, "topic_id", version).into());
629        }
630        write_compact_array_length(buf, self.partition_indexes.len() as i32);
631        for el in &self.partition_indexes {
632            write_i32(buf, *el);
633        }
634        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
635        all_tags.sort_by_key(|f| f.tag);
636        write_tagged_fields(buf, &all_tags)?;
637        Ok(())
638    }
639    pub fn encoded_len(&self, version: i16) -> Result<usize> {
640        let mut len: usize = 0;
641        if version <= 9 {
642            len += compact_string_len(&self.name)?;
643        } else if self.name != KafkaString::default() {
644            return Err(UnsupportedFieldVersion::new(9, "name", version).into());
645        }
646        if version >= 10 {
647            len += 16;
648        } else if self.topic_id != KafkaUuid::ZERO {
649            return Err(UnsupportedFieldVersion::new(9, "topic_id", version).into());
650        }
651        len += compact_array_length_len(self.partition_indexes.len() as i32);
652        len += self.partition_indexes.len() * 4usize;
653        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
654        all_tags.sort_by_key(|f| f.tag);
655        len += tagged_fields_len(&all_tags)?;
656        Ok(len)
657    }
658}