kacrab_protocol/generated/
join_group_response.rs1#![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 pub throttle_time_ms: i32,
20 pub error_code: i16,
22 pub generation_id: i32,
24 pub protocol_type: Option<KafkaString>,
26 pub protocol_name: Option<KafkaString>,
28 pub leader: KafkaString,
30 pub skip_assignment: bool,
32 pub member_id: KafkaString,
34 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 pub member_id: KafkaString,
305 pub group_instance_id: Option<KafkaString>,
307 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}