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 ProduceResponseData {
17 pub responses: Vec<TopicProduceResponse>,
19 pub throttle_time_ms: i32,
22 pub node_endpoints: Vec<NodeEndpoint>,
25 pub _unknown_tagged_fields: Vec<RawTaggedField>,
26}
27impl Default for ProduceResponseData {
28 fn default() -> Self {
29 Self {
30 responses: Vec::new(),
31 throttle_time_ms: 0i32,
32 node_endpoints: Vec::new(),
33 _unknown_tagged_fields: Vec::new(),
34 }
35 }
36}
37impl ProduceResponseData {
38 pub fn with_responses(mut self, value: Vec<TopicProduceResponse>) -> Self {
39 self.responses = value;
40 self
41 }
42 pub fn with_throttle_time_ms(mut self, value: i32) -> Self {
43 self.throttle_time_ms = value;
44 self
45 }
46 pub fn with_node_endpoints(mut self, value: Vec<NodeEndpoint>) -> Self {
47 self.node_endpoints = value;
48 self
49 }
50 pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
51 if version < 3 || version > 13 {
52 return Err(UnsupportedVersion::new(0, version).into());
53 }
54 let responses;
55 let throttle_time_ms;
56 let mut node_endpoints = Vec::new();
57 let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
58 if version >= 9 {
59 responses = {
60 let len = read_compact_array_length(buf)?;
61 let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
62 for _ in 0..len {
63 arr.push(TopicProduceResponse::read(buf, version)?);
64 }
65 arr
66 };
67 } else {
68 responses = {
69 let len = read_array_length(buf)?;
70 let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
71 for _ in 0..len {
72 arr.push(TopicProduceResponse::read(buf, version)?);
73 }
74 arr
75 };
76 }
77 throttle_time_ms = read_i32(buf)?;
78 if version >= 9 {
79 let tagged_fields = read_tagged_fields(buf)?;
80 for field in &tagged_fields {
81 match field.tag {
82 0 => {
83 if version >= 10 {
84 let mut tag_buf = field.data.clone();
85 node_endpoints = {
86 let len = read_compact_array_length(&mut tag_buf)?;
87 let mut arr = Vec::with_capacity(array_read_capacity(
88 len,
89 (&mut tag_buf).len(),
90 ));
91 for _ in 0..len {
92 arr.push(NodeEndpoint::read(&mut tag_buf, version)?);
93 }
94 arr
95 };
96 }
97 },
98 _ => {
99 _unknown_tagged_fields.push(field.clone());
100 },
101 }
102 }
103 }
104 Ok(Self {
105 responses,
106 throttle_time_ms,
107 node_endpoints,
108 _unknown_tagged_fields,
109 })
110 }
111 pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
112 if version < 3 || version > 13 {
113 return Err(UnsupportedVersion::new(0, version).into());
114 }
115 if version >= 9 {
116 write_compact_array_length(buf, self.responses.len() as i32);
117 for el in &self.responses {
118 el.write(buf, version)?;
119 }
120 } else {
121 write_array_length(buf, self.responses.len() as i32);
122 for el in &self.responses {
123 el.write(buf, version)?;
124 }
125 }
126 write_i32(buf, self.throttle_time_ms);
127 if version >= 9 {
128 let mut known_tagged_fields: Vec<RawTaggedField> = Vec::new();
129 if version >= 10 && !self.node_endpoints.is_empty() {
130 let mut tag_buf = BytesMut::new();
131 write_compact_array_length(&mut tag_buf, self.node_endpoints.len() as i32);
132 for el in &self.node_endpoints {
133 el.write(&mut tag_buf, version)?;
134 }
135 known_tagged_fields.push(RawTaggedField {
136 tag: 0,
137 data: tag_buf.freeze(),
138 });
139 }
140 let mut all_tags = known_tagged_fields;
141 all_tags.extend(self._unknown_tagged_fields.iter().cloned());
142 all_tags.sort_by_key(|f| f.tag);
143 write_tagged_fields(buf, &all_tags)?;
144 }
145 Ok(())
146 }
147 pub fn encoded_len(&self, version: i16) -> Result<usize> {
148 if version < 3 || version > 13 {
149 return Err(UnsupportedVersion::new(0, version).into());
150 }
151 let mut len: usize = 0;
152 if version >= 9 {
153 len += compact_array_length_len(self.responses.len() as i32);
154 for el in &self.responses {
155 len += el.encoded_len(version)?;
156 }
157 } else {
158 len += array_length_len();
159 for el in &self.responses {
160 len += el.encoded_len(version)?;
161 }
162 }
163 len += 4;
164 if version >= 9 {
165 let mut known_tagged_fields: Vec<RawTaggedField> = Vec::new();
166 if version >= 10 && !self.node_endpoints.is_empty() {
167 let mut tag_buf = BytesMut::new();
168 write_compact_array_length(&mut tag_buf, self.node_endpoints.len() as i32);
169 for el in &self.node_endpoints {
170 el.write(&mut tag_buf, version)?;
171 }
172 known_tagged_fields.push(RawTaggedField {
173 tag: 0,
174 data: tag_buf.freeze(),
175 });
176 }
177 let mut all_tags = known_tagged_fields;
178 all_tags.extend(self._unknown_tagged_fields.iter().cloned());
179 all_tags.sort_by_key(|f| f.tag);
180 len += tagged_fields_len(&all_tags)?;
181 }
182 Ok(len)
183 }
184}
185#[derive(Debug, Clone, PartialEq)]
186pub struct TopicProduceResponse {
187 pub name: KafkaString,
189 pub topic_id: KafkaUuid,
191 pub partition_responses: Vec<PartitionProduceResponse>,
193 pub _unknown_tagged_fields: Vec<RawTaggedField>,
194}
195impl Default for TopicProduceResponse {
196 fn default() -> Self {
197 Self {
198 name: KafkaString::default(),
199 topic_id: KafkaUuid::ZERO,
200 partition_responses: Vec::new(),
201 _unknown_tagged_fields: Vec::new(),
202 }
203 }
204}
205impl TopicProduceResponse {
206 pub fn with_name(mut self, value: KafkaString) -> Self {
207 self.name = value;
208 self
209 }
210 pub fn with_topic_id(mut self, value: KafkaUuid) -> Self {
211 self.topic_id = value;
212 self
213 }
214 pub fn with_partition_responses(mut self, value: Vec<PartitionProduceResponse>) -> Self {
215 self.partition_responses = value;
216 self
217 }
218 pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
219 let mut name = KafkaString::default();
220 let mut topic_id = KafkaUuid::ZERO;
221 let partition_responses;
222 let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
223 if version <= 12 {
224 if version >= 9 {
225 name = read_compact_string(buf)?;
226 } else {
227 name = read_string(buf)?;
228 }
229 }
230 if version >= 13 {
231 topic_id = read_uuid(buf)?;
232 }
233 if version >= 9 {
234 partition_responses = {
235 let len = read_compact_array_length(buf)?;
236 let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
237 for _ in 0..len {
238 arr.push(PartitionProduceResponse::read(buf, version)?);
239 }
240 arr
241 };
242 } else {
243 partition_responses = {
244 let len = read_array_length(buf)?;
245 let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
246 for _ in 0..len {
247 arr.push(PartitionProduceResponse::read(buf, version)?);
248 }
249 arr
250 };
251 }
252 if version >= 9 {
253 let tagged_fields = read_tagged_fields(buf)?;
254 for field in &tagged_fields {
255 match field.tag {
256 _ => {
257 _unknown_tagged_fields.push(field.clone());
258 },
259 }
260 }
261 }
262 Ok(Self {
263 name,
264 topic_id,
265 partition_responses,
266 _unknown_tagged_fields,
267 })
268 }
269 pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
270 if version <= 12 {
271 if version >= 9 {
272 write_compact_string(buf, &self.name)?;
273 } else {
274 write_string(buf, &self.name)?;
275 }
276 } else if self.name != KafkaString::default() {
277 return Err(UnsupportedFieldVersion::new(0, "name", version).into());
278 }
279 if version >= 13 {
280 write_uuid(buf, &self.topic_id);
281 } else if self.topic_id != KafkaUuid::ZERO {
282 return Err(UnsupportedFieldVersion::new(0, "topic_id", version).into());
283 }
284 if version >= 9 {
285 write_compact_array_length(buf, self.partition_responses.len() as i32);
286 for el in &self.partition_responses {
287 el.write(buf, version)?;
288 }
289 } else {
290 write_array_length(buf, self.partition_responses.len() as i32);
291 for el in &self.partition_responses {
292 el.write(buf, version)?;
293 }
294 }
295 if version >= 9 {
296 let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
297 all_tags.sort_by_key(|f| f.tag);
298 write_tagged_fields(buf, &all_tags)?;
299 }
300 Ok(())
301 }
302 pub fn encoded_len(&self, version: i16) -> Result<usize> {
303 let mut len: usize = 0;
304 if version <= 12 {
305 if version >= 9 {
306 len += compact_string_len(&self.name)?;
307 } else {
308 len += string_len(&self.name)?;
309 }
310 } else if self.name != KafkaString::default() {
311 return Err(UnsupportedFieldVersion::new(0, "name", version).into());
312 }
313 if version >= 13 {
314 len += 16;
315 } else if self.topic_id != KafkaUuid::ZERO {
316 return Err(UnsupportedFieldVersion::new(0, "topic_id", version).into());
317 }
318 if version >= 9 {
319 len += compact_array_length_len(self.partition_responses.len() as i32);
320 for el in &self.partition_responses {
321 len += el.encoded_len(version)?;
322 }
323 } else {
324 len += array_length_len();
325 for el in &self.partition_responses {
326 len += el.encoded_len(version)?;
327 }
328 }
329 if version >= 9 {
330 let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
331 all_tags.sort_by_key(|f| f.tag);
332 len += tagged_fields_len(&all_tags)?;
333 }
334 Ok(len)
335 }
336}
337#[derive(Debug, Clone, PartialEq)]
338pub struct PartitionProduceResponse {
339 pub index: i32,
341 pub error_code: i16,
343 pub base_offset: i64,
345 pub log_append_time_ms: i64,
349 pub log_start_offset: i64,
351 pub record_errors: Vec<BatchIndexAndErrorMessage>,
353 pub error_message: Option<KafkaString>,
356 pub current_leader: LeaderIdAndEpoch,
358 pub _unknown_tagged_fields: Vec<RawTaggedField>,
359}
360impl Default for PartitionProduceResponse {
361 fn default() -> Self {
362 Self {
363 index: 0_i32,
364 error_code: 0_i16,
365 base_offset: 0_i64,
366 log_append_time_ms: -1i64,
367 log_start_offset: -1i64,
368 record_errors: Vec::new(),
369 error_message: None,
370 current_leader: LeaderIdAndEpoch::default(),
371 _unknown_tagged_fields: Vec::new(),
372 }
373 }
374}
375impl PartitionProduceResponse {
376 pub fn with_index(mut self, value: i32) -> Self {
377 self.index = value;
378 self
379 }
380 pub fn with_error_code(mut self, value: i16) -> Self {
381 self.error_code = value;
382 self
383 }
384 pub fn with_base_offset(mut self, value: i64) -> Self {
385 self.base_offset = value;
386 self
387 }
388 pub fn with_log_append_time_ms(mut self, value: i64) -> Self {
389 self.log_append_time_ms = value;
390 self
391 }
392 pub fn with_log_start_offset(mut self, value: i64) -> Self {
393 self.log_start_offset = value;
394 self
395 }
396 pub fn with_record_errors(mut self, value: Vec<BatchIndexAndErrorMessage>) -> Self {
397 self.record_errors = value;
398 self
399 }
400 pub fn with_error_message(mut self, value: Option<KafkaString>) -> Self {
401 self.error_message = value;
402 self
403 }
404 pub fn with_current_leader(mut self, value: LeaderIdAndEpoch) -> Self {
405 self.current_leader = value;
406 self
407 }
408 pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
409 let index;
410 let error_code;
411 let base_offset;
412 let log_append_time_ms;
413 let mut log_start_offset = -1i64;
414 let mut record_errors = Vec::new();
415 let mut error_message = None;
416 let mut current_leader = LeaderIdAndEpoch::default();
417 let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
418 index = read_i32(buf)?;
419 error_code = read_i16(buf)?;
420 base_offset = read_i64(buf)?;
421 log_append_time_ms = read_i64(buf)?;
422 if version >= 5 {
423 log_start_offset = read_i64(buf)?;
424 }
425 if version >= 8 {
426 if version >= 9 {
427 record_errors = {
428 let len = read_compact_array_length(buf)?;
429 let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
430 for _ in 0..len {
431 arr.push(BatchIndexAndErrorMessage::read(buf, version)?);
432 }
433 arr
434 };
435 } else {
436 record_errors = {
437 let len = read_array_length(buf)?;
438 let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
439 for _ in 0..len {
440 arr.push(BatchIndexAndErrorMessage::read(buf, version)?);
441 }
442 arr
443 };
444 }
445 }
446 if version >= 8 {
447 if version >= 9 {
448 error_message = read_compact_nullable_string(buf)?;
449 } else {
450 error_message = read_nullable_string(buf)?;
451 }
452 }
453 if version >= 9 {
454 let tagged_fields = read_tagged_fields(buf)?;
455 for field in &tagged_fields {
456 match field.tag {
457 0 => {
458 if version >= 10 {
459 let mut tag_buf = field.data.clone();
460 current_leader = LeaderIdAndEpoch::read(&mut tag_buf, version)?;
461 }
462 },
463 _ => {
464 _unknown_tagged_fields.push(field.clone());
465 },
466 }
467 }
468 }
469 Ok(Self {
470 index,
471 error_code,
472 base_offset,
473 log_append_time_ms,
474 log_start_offset,
475 record_errors,
476 error_message,
477 current_leader,
478 _unknown_tagged_fields,
479 })
480 }
481 pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
482 write_i32(buf, self.index);
483 write_i16(buf, self.error_code);
484 write_i64(buf, self.base_offset);
485 write_i64(buf, self.log_append_time_ms);
486 if version >= 5 {
487 write_i64(buf, self.log_start_offset);
488 } else if self.log_start_offset != -1i64 {
489 return Err(UnsupportedFieldVersion::new(0, "log_start_offset", version).into());
490 }
491 if version >= 8 {
492 if version >= 9 {
493 write_compact_array_length(buf, self.record_errors.len() as i32);
494 for el in &self.record_errors {
495 el.write(buf, version)?;
496 }
497 } else {
498 write_array_length(buf, self.record_errors.len() as i32);
499 for el in &self.record_errors {
500 el.write(buf, version)?;
501 }
502 }
503 } else if self.record_errors != Vec::new() {
504 return Err(UnsupportedFieldVersion::new(0, "record_errors", version).into());
505 }
506 if version >= 8 {
507 if version >= 9 {
508 write_compact_nullable_string(buf, self.error_message.as_ref())?;
509 } else {
510 write_nullable_string(buf, self.error_message.as_ref())?;
511 }
512 } else if self.error_message != None {
513 return Err(UnsupportedFieldVersion::new(0, "error_message", version).into());
514 }
515 if version >= 9 {
516 let mut known_tagged_fields: Vec<RawTaggedField> = Vec::new();
517 if version >= 10 && self.current_leader != LeaderIdAndEpoch::default() {
518 let mut tag_buf = BytesMut::new();
519 self.current_leader.write(&mut tag_buf, version)?;
520 known_tagged_fields.push(RawTaggedField {
521 tag: 0,
522 data: tag_buf.freeze(),
523 });
524 }
525 let mut all_tags = known_tagged_fields;
526 all_tags.extend(self._unknown_tagged_fields.iter().cloned());
527 all_tags.sort_by_key(|f| f.tag);
528 write_tagged_fields(buf, &all_tags)?;
529 }
530 Ok(())
531 }
532 pub fn encoded_len(&self, version: i16) -> Result<usize> {
533 let mut len: usize = 0;
534 len += 4;
535 len += 2;
536 len += 8;
537 len += 8;
538 if version >= 5 {
539 len += 8;
540 } else if self.log_start_offset != -1i64 {
541 return Err(UnsupportedFieldVersion::new(0, "log_start_offset", version).into());
542 }
543 if version >= 8 {
544 if version >= 9 {
545 len += compact_array_length_len(self.record_errors.len() as i32);
546 for el in &self.record_errors {
547 len += el.encoded_len(version)?;
548 }
549 } else {
550 len += array_length_len();
551 for el in &self.record_errors {
552 len += el.encoded_len(version)?;
553 }
554 }
555 } else if self.record_errors != Vec::new() {
556 return Err(UnsupportedFieldVersion::new(0, "record_errors", version).into());
557 }
558 if version >= 8 {
559 if version >= 9 {
560 len += compact_nullable_string_len(self.error_message.as_ref())?;
561 } else {
562 len += nullable_string_len(self.error_message.as_ref())?;
563 }
564 } else if self.error_message != None {
565 return Err(UnsupportedFieldVersion::new(0, "error_message", version).into());
566 }
567 if version >= 9 {
568 let mut known_tagged_fields: Vec<RawTaggedField> = Vec::new();
569 if version >= 10 && self.current_leader != LeaderIdAndEpoch::default() {
570 let mut tag_buf = BytesMut::new();
571 self.current_leader.write(&mut tag_buf, version)?;
572 known_tagged_fields.push(RawTaggedField {
573 tag: 0,
574 data: tag_buf.freeze(),
575 });
576 }
577 let mut all_tags = known_tagged_fields;
578 all_tags.extend(self._unknown_tagged_fields.iter().cloned());
579 all_tags.sort_by_key(|f| f.tag);
580 len += tagged_fields_len(&all_tags)?;
581 }
582 Ok(len)
583 }
584}
585#[derive(Debug, Clone, PartialEq)]
586pub struct BatchIndexAndErrorMessage {
587 pub batch_index: i32,
589 pub batch_index_error_message: Option<KafkaString>,
591 pub _unknown_tagged_fields: Vec<RawTaggedField>,
592}
593impl Default for BatchIndexAndErrorMessage {
594 fn default() -> Self {
595 Self {
596 batch_index: 0_i32,
597 batch_index_error_message: None,
598 _unknown_tagged_fields: Vec::new(),
599 }
600 }
601}
602impl BatchIndexAndErrorMessage {
603 pub fn with_batch_index(mut self, value: i32) -> Self {
604 self.batch_index = value;
605 self
606 }
607 pub fn with_batch_index_error_message(mut self, value: Option<KafkaString>) -> Self {
608 self.batch_index_error_message = value;
609 self
610 }
611 pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
612 let batch_index;
613 let batch_index_error_message;
614 let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
615 batch_index = read_i32(buf)?;
616 if version >= 9 {
617 batch_index_error_message = read_compact_nullable_string(buf)?;
618 } else {
619 batch_index_error_message = read_nullable_string(buf)?;
620 }
621 if version >= 9 {
622 let tagged_fields = read_tagged_fields(buf)?;
623 for field in &tagged_fields {
624 match field.tag {
625 _ => {
626 _unknown_tagged_fields.push(field.clone());
627 },
628 }
629 }
630 }
631 Ok(Self {
632 batch_index,
633 batch_index_error_message,
634 _unknown_tagged_fields,
635 })
636 }
637 pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
638 write_i32(buf, self.batch_index);
639 if version >= 9 {
640 write_compact_nullable_string(buf, self.batch_index_error_message.as_ref())?;
641 } else {
642 write_nullable_string(buf, self.batch_index_error_message.as_ref())?;
643 }
644 if version >= 9 {
645 let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
646 all_tags.sort_by_key(|f| f.tag);
647 write_tagged_fields(buf, &all_tags)?;
648 }
649 Ok(())
650 }
651 pub fn encoded_len(&self, version: i16) -> Result<usize> {
652 let mut len: usize = 0;
653 len += 4;
654 if version >= 9 {
655 len += compact_nullable_string_len(self.batch_index_error_message.as_ref())?;
656 } else {
657 len += nullable_string_len(self.batch_index_error_message.as_ref())?;
658 }
659 if version >= 9 {
660 let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
661 all_tags.sort_by_key(|f| f.tag);
662 len += tagged_fields_len(&all_tags)?;
663 }
664 Ok(len)
665 }
666}
667#[derive(Debug, Clone, PartialEq)]
668pub struct LeaderIdAndEpoch {
669 pub leader_id: i32,
671 pub leader_epoch: i32,
673 pub _unknown_tagged_fields: Vec<RawTaggedField>,
674}
675impl Default for LeaderIdAndEpoch {
676 fn default() -> Self {
677 Self {
678 leader_id: -1i32,
679 leader_epoch: -1i32,
680 _unknown_tagged_fields: Vec::new(),
681 }
682 }
683}
684impl LeaderIdAndEpoch {
685 pub fn with_leader_id(mut self, value: i32) -> Self {
686 self.leader_id = value;
687 self
688 }
689 pub fn with_leader_epoch(mut self, value: i32) -> Self {
690 self.leader_epoch = value;
691 self
692 }
693 pub fn read(buf: &mut Bytes, _version: i16) -> Result<Self> {
694 let leader_id;
695 let leader_epoch;
696 let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
697 leader_id = read_i32(buf)?;
698 leader_epoch = read_i32(buf)?;
699 let tagged_fields = read_tagged_fields(buf)?;
700 for field in &tagged_fields {
701 match field.tag {
702 _ => {
703 _unknown_tagged_fields.push(field.clone());
704 },
705 }
706 }
707 Ok(Self {
708 leader_id,
709 leader_epoch,
710 _unknown_tagged_fields,
711 })
712 }
713 pub fn write(&self, buf: &mut BytesMut, _version: i16) -> Result<()> {
714 write_i32(buf, self.leader_id);
715 write_i32(buf, self.leader_epoch);
716 let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
717 all_tags.sort_by_key(|f| f.tag);
718 write_tagged_fields(buf, &all_tags)?;
719 Ok(())
720 }
721 pub fn encoded_len(&self, _version: i16) -> Result<usize> {
722 let mut len: usize = 0;
723 len += 4;
724 len += 4;
725 let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
726 all_tags.sort_by_key(|f| f.tag);
727 len += tagged_fields_len(&all_tags)?;
728 Ok(len)
729 }
730}
731#[derive(Debug, Clone, PartialEq)]
732pub struct NodeEndpoint {
733 pub node_id: i32,
735 pub host: KafkaString,
737 pub port: i32,
739 pub rack: Option<KafkaString>,
741 pub _unknown_tagged_fields: Vec<RawTaggedField>,
742}
743impl Default for NodeEndpoint {
744 fn default() -> Self {
745 Self {
746 node_id: 0_i32,
747 host: KafkaString::default(),
748 port: 0_i32,
749 rack: None,
750 _unknown_tagged_fields: Vec::new(),
751 }
752 }
753}
754impl NodeEndpoint {
755 pub fn with_node_id(mut self, value: i32) -> Self {
756 self.node_id = value;
757 self
758 }
759 pub fn with_host(mut self, value: KafkaString) -> Self {
760 self.host = value;
761 self
762 }
763 pub fn with_port(mut self, value: i32) -> Self {
764 self.port = value;
765 self
766 }
767 pub fn with_rack(mut self, value: Option<KafkaString>) -> Self {
768 self.rack = value;
769 self
770 }
771 pub fn read(buf: &mut Bytes, _version: i16) -> Result<Self> {
772 let node_id;
773 let host;
774 let port;
775 let rack;
776 let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
777 node_id = read_i32(buf)?;
778 host = read_compact_string(buf)?;
779 port = read_i32(buf)?;
780 rack = read_compact_nullable_string(buf)?;
781 let tagged_fields = read_tagged_fields(buf)?;
782 for field in &tagged_fields {
783 match field.tag {
784 _ => {
785 _unknown_tagged_fields.push(field.clone());
786 },
787 }
788 }
789 Ok(Self {
790 node_id,
791 host,
792 port,
793 rack,
794 _unknown_tagged_fields,
795 })
796 }
797 pub fn write(&self, buf: &mut BytesMut, _version: i16) -> Result<()> {
798 write_i32(buf, self.node_id);
799 write_compact_string(buf, &self.host)?;
800 write_i32(buf, self.port);
801 write_compact_nullable_string(buf, self.rack.as_ref())?;
802 let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
803 all_tags.sort_by_key(|f| f.tag);
804 write_tagged_fields(buf, &all_tags)?;
805 Ok(())
806 }
807 pub fn encoded_len(&self, _version: i16) -> Result<usize> {
808 let mut len: usize = 0;
809 len += 4;
810 len += compact_string_len(&self.host)?;
811 len += 4;
812 len += compact_nullable_string_len(self.rack.as_ref())?;
813 let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
814 all_tags.sort_by_key(|f| f.tag);
815 len += tagged_fields_len(&all_tags)?;
816 Ok(len)
817 }
818}