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 DescribeQuorumResponseData {
17 pub error_code: i16,
19 pub error_message: Option<KafkaString>,
21 pub topics: Vec<TopicData>,
23 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 pub topic_name: KafkaString,
162 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 pub partition_index: i32,
239 pub error_code: i16,
241 pub error_message: Option<KafkaString>,
243 pub leader_id: i32,
245 pub leader_epoch: i32,
247 pub high_watermark: i64,
249 pub current_voters: Vec<ReplicaState>,
251 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 pub node_id: i32,
411 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 pub name: KafkaString,
488 pub host: KafkaString,
490 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 pub replica_id: i32,
564 pub replica_directory_id: KafkaUuid,
566 pub log_end_offset: i64,
568 pub last_fetch_timestamp: i64,
571 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}