1#![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 pub group_id: KafkaString,
19 pub topics: Option<Vec<OffsetFetchRequestTopic>>,
21 pub groups: Vec<OffsetFetchRequestGroup>,
23 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(array_read_capacity(len, (buf).len()));
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(array_read_capacity(len, (buf).len()));
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(array_read_capacity(len, (buf).len()));
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(array_read_capacity(len, (buf).len()));
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 pub name: KafkaString,
301 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(array_read_capacity(len, (buf).len()));
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(array_read_capacity(len, (buf).len()));
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 pub group_id: KafkaString,
417 pub member_id: Option<KafkaString>,
419 pub member_epoch: i32,
421 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(array_read_capacity(len, (buf).len()));
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 pub name: KafkaString,
556 pub topic_id: KafkaUuid,
558 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(array_read_capacity(len, (buf).len()));
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}