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