Skip to main content

kacrab_protocol/generated/
join_group_response.rs

1//! Generated from JoinGroupResponse.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 JoinGroupResponseData {
17    /// The duration in milliseconds for which the request was throttled due to a quota violation,
18    /// or zero if the request did not violate any quota.
19    pub throttle_time_ms: i32,
20    /// The error code, or 0 if there was no error.
21    pub error_code: i16,
22    /// The generation ID of the group.
23    pub generation_id: i32,
24    /// The group protocol name.
25    pub protocol_type: Option<KafkaString>,
26    /// The group protocol selected by the coordinator.
27    pub protocol_name: Option<KafkaString>,
28    /// The leader of the group.
29    pub leader: KafkaString,
30    /// True if the leader must skip running the assignment.
31    pub skip_assignment: bool,
32    /// The member ID assigned by the group coordinator.
33    pub member_id: KafkaString,
34    /// The group members.
35    pub members: Vec<JoinGroupResponseMember>,
36    pub _unknown_tagged_fields: Vec<RawTaggedField>,
37}
38impl Default for JoinGroupResponseData {
39    fn default() -> Self {
40        Self {
41            throttle_time_ms: 0_i32,
42            error_code: 0_i16,
43            generation_id: -1i32,
44            protocol_type: None,
45            protocol_name: None,
46            leader: KafkaString::default(),
47            skip_assignment: false,
48            member_id: KafkaString::default(),
49            members: Vec::new(),
50            _unknown_tagged_fields: Vec::new(),
51        }
52    }
53}
54impl JoinGroupResponseData {
55    pub fn with_throttle_time_ms(mut self, value: i32) -> Self {
56        self.throttle_time_ms = value;
57        self
58    }
59    pub fn with_error_code(mut self, value: i16) -> Self {
60        self.error_code = value;
61        self
62    }
63    pub fn with_generation_id(mut self, value: i32) -> Self {
64        self.generation_id = value;
65        self
66    }
67    pub fn with_protocol_type(mut self, value: Option<KafkaString>) -> Self {
68        self.protocol_type = value;
69        self
70    }
71    pub fn with_protocol_name(mut self, value: Option<KafkaString>) -> Self {
72        self.protocol_name = value;
73        self
74    }
75    pub fn with_leader(mut self, value: KafkaString) -> Self {
76        self.leader = value;
77        self
78    }
79    pub fn with_skip_assignment(mut self, value: bool) -> Self {
80        self.skip_assignment = value;
81        self
82    }
83    pub fn with_member_id(mut self, value: KafkaString) -> Self {
84        self.member_id = value;
85        self
86    }
87    pub fn with_members(mut self, value: Vec<JoinGroupResponseMember>) -> Self {
88        self.members = value;
89        self
90    }
91    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
92        if version < 0 || version > 9 {
93            return Err(UnsupportedVersion::new(11, version).into());
94        }
95        let mut throttle_time_ms = 0_i32;
96        let error_code;
97        let generation_id;
98        let mut protocol_type = None;
99        let protocol_name;
100        let leader;
101        let mut skip_assignment = false;
102        let member_id;
103        let members;
104        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
105        if version >= 2 {
106            throttle_time_ms = read_i32(buf)?;
107        }
108        error_code = read_i16(buf)?;
109        generation_id = read_i32(buf)?;
110        if version >= 7 {
111            protocol_type = read_compact_nullable_string(buf)?;
112        }
113        if version >= 7 {
114            protocol_name = read_compact_nullable_string(buf)?;
115        } else {
116            if version >= 6 {
117                protocol_name = Some(read_compact_string(buf)?);
118            } else {
119                protocol_name = Some(read_string(buf)?);
120            }
121        }
122        if version >= 6 {
123            leader = read_compact_string(buf)?;
124        } else {
125            leader = read_string(buf)?;
126        }
127        if version >= 9 {
128            skip_assignment = read_bool(buf)?;
129        }
130        if version >= 6 {
131            member_id = read_compact_string(buf)?;
132        } else {
133            member_id = read_string(buf)?;
134        }
135        if version >= 6 {
136            members = {
137                let len = read_compact_array_length(buf)?;
138                let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
139                for _ in 0..len {
140                    arr.push(JoinGroupResponseMember::read(buf, version)?);
141                }
142                arr
143            };
144        } else {
145            members = {
146                let len = read_array_length(buf)?;
147                let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
148                for _ in 0..len {
149                    arr.push(JoinGroupResponseMember::read(buf, version)?);
150                }
151                arr
152            };
153        }
154        if version >= 6 {
155            let tagged_fields = read_tagged_fields(buf)?;
156            for field in &tagged_fields {
157                match field.tag {
158                    _ => {
159                        _unknown_tagged_fields.push(field.clone());
160                    },
161                }
162            }
163        }
164        Ok(Self {
165            throttle_time_ms,
166            error_code,
167            generation_id,
168            protocol_type,
169            protocol_name,
170            leader,
171            skip_assignment,
172            member_id,
173            members,
174            _unknown_tagged_fields,
175        })
176    }
177    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
178        if version < 0 || version > 9 {
179            return Err(UnsupportedVersion::new(11, version).into());
180        }
181        if version >= 2 {
182            write_i32(buf, self.throttle_time_ms);
183        } else if self.throttle_time_ms != 0_i32 {
184            return Err(UnsupportedFieldVersion::new(11, "throttle_time_ms", version).into());
185        }
186        write_i16(buf, self.error_code);
187        write_i32(buf, self.generation_id);
188        if version >= 7 {
189            write_compact_nullable_string(buf, self.protocol_type.as_ref())?;
190        } else if self.protocol_type != None {
191            return Err(UnsupportedFieldVersion::new(11, "protocol_type", version).into());
192        }
193        if version >= 7 {
194            write_compact_nullable_string(buf, self.protocol_name.as_ref())?;
195        } else {
196            {
197                let _nn_default = KafkaString::default();
198                let _nn_val = self.protocol_name.as_ref().unwrap_or(&_nn_default);
199                if version >= 6 {
200                    write_compact_string(buf, _nn_val)?;
201                } else {
202                    write_string(buf, _nn_val)?;
203                }
204            }
205        }
206        if version >= 6 {
207            write_compact_string(buf, &self.leader)?;
208        } else {
209            write_string(buf, &self.leader)?;
210        }
211        if version >= 9 {
212            write_bool(buf, self.skip_assignment);
213        } else if self.skip_assignment != false {
214            return Err(UnsupportedFieldVersion::new(11, "skip_assignment", version).into());
215        }
216        if version >= 6 {
217            write_compact_string(buf, &self.member_id)?;
218        } else {
219            write_string(buf, &self.member_id)?;
220        }
221        if version >= 6 {
222            write_compact_array_length(buf, self.members.len() as i32);
223            for el in &self.members {
224                el.write(buf, version)?;
225            }
226        } else {
227            write_array_length(buf, self.members.len() as i32);
228            for el in &self.members {
229                el.write(buf, version)?;
230            }
231        }
232        if version >= 6 {
233            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
234            all_tags.sort_by_key(|f| f.tag);
235            write_tagged_fields(buf, &all_tags)?;
236        }
237        Ok(())
238    }
239    pub fn encoded_len(&self, version: i16) -> Result<usize> {
240        if version < 0 || version > 9 {
241            return Err(UnsupportedVersion::new(11, version).into());
242        }
243        let mut len: usize = 0;
244        if version >= 2 {
245            len += 4;
246        } else if self.throttle_time_ms != 0_i32 {
247            return Err(UnsupportedFieldVersion::new(11, "throttle_time_ms", version).into());
248        }
249        len += 2;
250        len += 4;
251        if version >= 7 {
252            len += compact_nullable_string_len(self.protocol_type.as_ref())?;
253        } else if self.protocol_type != None {
254            return Err(UnsupportedFieldVersion::new(11, "protocol_type", version).into());
255        }
256        if version >= 7 {
257            len += compact_nullable_string_len(self.protocol_name.as_ref())?;
258        } else {
259            let _nn_default = KafkaString::default();
260            let _nn_val = self.protocol_name.as_ref().unwrap_or(&_nn_default);
261            if version >= 6 {
262                len += compact_string_len(_nn_val)?;
263            } else {
264                len += string_len(_nn_val)?;
265            }
266        }
267        if version >= 6 {
268            len += compact_string_len(&self.leader)?;
269        } else {
270            len += string_len(&self.leader)?;
271        }
272        if version >= 9 {
273            len += 1;
274        } else if self.skip_assignment != false {
275            return Err(UnsupportedFieldVersion::new(11, "skip_assignment", version).into());
276        }
277        if version >= 6 {
278            len += compact_string_len(&self.member_id)?;
279        } else {
280            len += string_len(&self.member_id)?;
281        }
282        if version >= 6 {
283            len += compact_array_length_len(self.members.len() as i32);
284            for el in &self.members {
285                len += el.encoded_len(version)?;
286            }
287        } else {
288            len += array_length_len();
289            for el in &self.members {
290                len += el.encoded_len(version)?;
291            }
292        }
293        if version >= 6 {
294            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
295            all_tags.sort_by_key(|f| f.tag);
296            len += tagged_fields_len(&all_tags)?;
297        }
298        Ok(len)
299    }
300}
301#[derive(Debug, Clone, PartialEq)]
302pub struct JoinGroupResponseMember {
303    /// The group member ID.
304    pub member_id: KafkaString,
305    /// The unique identifier of the consumer instance provided by end user.
306    pub group_instance_id: Option<KafkaString>,
307    /// The group member metadata.
308    pub metadata: Bytes,
309    pub _unknown_tagged_fields: Vec<RawTaggedField>,
310}
311impl Default for JoinGroupResponseMember {
312    fn default() -> Self {
313        Self {
314            member_id: KafkaString::default(),
315            group_instance_id: None,
316            metadata: Bytes::new(),
317            _unknown_tagged_fields: Vec::new(),
318        }
319    }
320}
321impl JoinGroupResponseMember {
322    pub fn with_member_id(mut self, value: KafkaString) -> Self {
323        self.member_id = value;
324        self
325    }
326    pub fn with_group_instance_id(mut self, value: Option<KafkaString>) -> Self {
327        self.group_instance_id = value;
328        self
329    }
330    pub fn with_metadata(mut self, value: Bytes) -> Self {
331        self.metadata = value;
332        self
333    }
334    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
335        let member_id;
336        let mut group_instance_id = None;
337        let metadata;
338        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
339        if version >= 6 {
340            member_id = read_compact_string(buf)?;
341        } else {
342            member_id = read_string(buf)?;
343        }
344        if version >= 5 {
345            if version >= 6 {
346                group_instance_id = read_compact_nullable_string(buf)?;
347            } else {
348                group_instance_id = read_nullable_string(buf)?;
349            }
350        }
351        if version >= 6 {
352            metadata = read_compact_bytes(buf)?;
353        } else {
354            metadata = read_bytes(buf)?;
355        }
356        if version >= 6 {
357            let tagged_fields = read_tagged_fields(buf)?;
358            for field in &tagged_fields {
359                match field.tag {
360                    _ => {
361                        _unknown_tagged_fields.push(field.clone());
362                    },
363                }
364            }
365        }
366        Ok(Self {
367            member_id,
368            group_instance_id,
369            metadata,
370            _unknown_tagged_fields,
371        })
372    }
373    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
374        if version >= 6 {
375            write_compact_string(buf, &self.member_id)?;
376        } else {
377            write_string(buf, &self.member_id)?;
378        }
379        if version >= 5 {
380            if version >= 6 {
381                write_compact_nullable_string(buf, self.group_instance_id.as_ref())?;
382            } else {
383                write_nullable_string(buf, self.group_instance_id.as_ref())?;
384            }
385        } else if self.group_instance_id != None {
386            return Err(UnsupportedFieldVersion::new(11, "group_instance_id", version).into());
387        }
388        if version >= 6 {
389            write_compact_bytes(buf, &self.metadata)?;
390        } else {
391            write_bytes(buf, &self.metadata)?;
392        }
393        if version >= 6 {
394            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
395            all_tags.sort_by_key(|f| f.tag);
396            write_tagged_fields(buf, &all_tags)?;
397        }
398        Ok(())
399    }
400    pub fn encoded_len(&self, version: i16) -> Result<usize> {
401        let mut len: usize = 0;
402        if version >= 6 {
403            len += compact_string_len(&self.member_id)?;
404        } else {
405            len += string_len(&self.member_id)?;
406        }
407        if version >= 5 {
408            if version >= 6 {
409                len += compact_nullable_string_len(self.group_instance_id.as_ref())?;
410            } else {
411                len += nullable_string_len(self.group_instance_id.as_ref())?;
412            }
413        } else if self.group_instance_id != None {
414            return Err(UnsupportedFieldVersion::new(11, "group_instance_id", version).into());
415        }
416        if version >= 6 {
417            len += compact_bytes_len(&self.metadata)?;
418        } else {
419            len += bytes_len(&self.metadata)?;
420        }
421        if version >= 6 {
422            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
423            all_tags.sort_by_key(|f| f.tag);
424            len += tagged_fields_len(&all_tags)?;
425        }
426        Ok(len)
427    }
428}