Skip to main content

kacrab_protocol/generated/
share_fetch_response.rs

1//! Generated from ShareFetchResponse.json - DO NOT EDIT
2#![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 ShareFetchResponseData {
17    /// The duration in milliseconds for which the request was throttled due to a quota violation,
18    /// or zero if the request did not violate any quota.
19    pub throttle_time_ms: i32,
20    /// The top-level response error code.
21    pub error_code: i16,
22    /// The top-level error message, or null if there was no error.
23    pub error_message: Option<KafkaString>,
24    /// The time in milliseconds for which the acquired records are locked.
25    pub acquisition_lock_timeout_ms: i32,
26    /// The response topics.
27    pub responses: Vec<ShareFetchableTopicResponse>,
28    /// Endpoints for all current leaders enumerated in PartitionData with error
29    /// NOT_LEADER_OR_FOLLOWER.
30    pub node_endpoints: Vec<NodeEndpoint>,
31    pub _unknown_tagged_fields: Vec<RawTaggedField>,
32}
33impl Default for ShareFetchResponseData {
34    fn default() -> Self {
35        Self {
36            throttle_time_ms: 0_i32,
37            error_code: 0_i16,
38            error_message: None,
39            acquisition_lock_timeout_ms: 0_i32,
40            responses: Vec::new(),
41            node_endpoints: Vec::new(),
42            _unknown_tagged_fields: Vec::new(),
43        }
44    }
45}
46impl ShareFetchResponseData {
47    pub fn with_throttle_time_ms(mut self, value: i32) -> Self {
48        self.throttle_time_ms = value;
49        self
50    }
51    pub fn with_error_code(mut self, value: i16) -> Self {
52        self.error_code = value;
53        self
54    }
55    pub fn with_error_message(mut self, value: Option<KafkaString>) -> Self {
56        self.error_message = value;
57        self
58    }
59    pub fn with_acquisition_lock_timeout_ms(mut self, value: i32) -> Self {
60        self.acquisition_lock_timeout_ms = value;
61        self
62    }
63    pub fn with_responses(mut self, value: Vec<ShareFetchableTopicResponse>) -> Self {
64        self.responses = value;
65        self
66    }
67    pub fn with_node_endpoints(mut self, value: Vec<NodeEndpoint>) -> Self {
68        self.node_endpoints = value;
69        self
70    }
71    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
72        if version < 1 || version > 2 {
73            return Err(UnsupportedVersion::new(78, version).into());
74        }
75        let throttle_time_ms;
76        let error_code;
77        let error_message;
78        let acquisition_lock_timeout_ms;
79        let responses;
80        let node_endpoints;
81        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
82        throttle_time_ms = read_i32(buf)?;
83        error_code = read_i16(buf)?;
84        error_message = read_compact_nullable_string(buf)?;
85        acquisition_lock_timeout_ms = read_i32(buf)?;
86        responses = {
87            let len = read_compact_array_length(buf)?;
88            let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
89            for _ in 0..len {
90                arr.push(ShareFetchableTopicResponse::read(buf, version)?);
91            }
92            arr
93        };
94        node_endpoints = {
95            let len = read_compact_array_length(buf)?;
96            let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
97            for _ in 0..len {
98                arr.push(NodeEndpoint::read(buf, version)?);
99            }
100            arr
101        };
102        let tagged_fields = read_tagged_fields(buf)?;
103        for field in &tagged_fields {
104            match field.tag {
105                _ => {
106                    _unknown_tagged_fields.push(field.clone());
107                },
108            }
109        }
110        Ok(Self {
111            throttle_time_ms,
112            error_code,
113            error_message,
114            acquisition_lock_timeout_ms,
115            responses,
116            node_endpoints,
117            _unknown_tagged_fields,
118        })
119    }
120    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
121        if version < 1 || version > 2 {
122            return Err(UnsupportedVersion::new(78, version).into());
123        }
124        write_i32(buf, self.throttle_time_ms);
125        write_i16(buf, self.error_code);
126        write_compact_nullable_string(buf, self.error_message.as_ref())?;
127        write_i32(buf, self.acquisition_lock_timeout_ms);
128        write_compact_array_length(buf, self.responses.len() as i32);
129        for el in &self.responses {
130            el.write(buf, version)?;
131        }
132        write_compact_array_length(buf, self.node_endpoints.len() as i32);
133        for el in &self.node_endpoints {
134            el.write(buf, version)?;
135        }
136        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
137        all_tags.sort_by_key(|f| f.tag);
138        write_tagged_fields(buf, &all_tags)?;
139        Ok(())
140    }
141    pub fn encoded_len(&self, version: i16) -> Result<usize> {
142        if version < 1 || version > 2 {
143            return Err(UnsupportedVersion::new(78, version).into());
144        }
145        let mut len: usize = 0;
146        len += 4;
147        len += 2;
148        len += compact_nullable_string_len(self.error_message.as_ref())?;
149        len += 4;
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        len += compact_array_length_len(self.node_endpoints.len() as i32);
155        for el in &self.node_endpoints {
156            len += el.encoded_len(version)?;
157        }
158        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
159        all_tags.sort_by_key(|f| f.tag);
160        len += tagged_fields_len(&all_tags)?;
161        Ok(len)
162    }
163}
164#[derive(Debug, Clone, PartialEq)]
165pub struct ShareFetchableTopicResponse {
166    /// The unique topic ID.
167    pub topic_id: KafkaUuid,
168    /// The topic partitions.
169    pub partitions: Vec<PartitionData>,
170    pub _unknown_tagged_fields: Vec<RawTaggedField>,
171}
172impl Default for ShareFetchableTopicResponse {
173    fn default() -> Self {
174        Self {
175            topic_id: KafkaUuid::ZERO,
176            partitions: Vec::new(),
177            _unknown_tagged_fields: Vec::new(),
178        }
179    }
180}
181impl ShareFetchableTopicResponse {
182    pub fn with_topic_id(mut self, value: KafkaUuid) -> Self {
183        self.topic_id = value;
184        self
185    }
186    pub fn with_partitions(mut self, value: Vec<PartitionData>) -> Self {
187        self.partitions = value;
188        self
189    }
190    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
191        let topic_id;
192        let partitions;
193        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
194        topic_id = read_uuid(buf)?;
195        partitions = {
196            let len = read_compact_array_length(buf)?;
197            let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
198            for _ in 0..len {
199                arr.push(PartitionData::read(buf, version)?);
200            }
201            arr
202        };
203        let tagged_fields = read_tagged_fields(buf)?;
204        for field in &tagged_fields {
205            match field.tag {
206                _ => {
207                    _unknown_tagged_fields.push(field.clone());
208                },
209            }
210        }
211        Ok(Self {
212            topic_id,
213            partitions,
214            _unknown_tagged_fields,
215        })
216    }
217    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
218        write_uuid(buf, &self.topic_id);
219        write_compact_array_length(buf, self.partitions.len() as i32);
220        for el in &self.partitions {
221            el.write(buf, version)?;
222        }
223        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
224        all_tags.sort_by_key(|f| f.tag);
225        write_tagged_fields(buf, &all_tags)?;
226        Ok(())
227    }
228    pub fn encoded_len(&self, version: i16) -> Result<usize> {
229        let mut len: usize = 0;
230        len += 16;
231        len += compact_array_length_len(self.partitions.len() as i32);
232        for el in &self.partitions {
233            len += el.encoded_len(version)?;
234        }
235        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
236        all_tags.sort_by_key(|f| f.tag);
237        len += tagged_fields_len(&all_tags)?;
238        Ok(len)
239    }
240}
241#[derive(Debug, Clone, PartialEq)]
242pub struct PartitionData {
243    /// The partition index.
244    pub partition_index: i32,
245    /// The fetch error code, or 0 if there was no fetch error.
246    pub error_code: i16,
247    /// The fetch error message, or null if there was no fetch error.
248    pub error_message: Option<KafkaString>,
249    /// The acknowledge error code, or 0 if there was no acknowledge error.
250    pub acknowledge_error_code: i16,
251    /// The acknowledge error message, or null if there was no acknowledge error.
252    pub acknowledge_error_message: Option<KafkaString>,
253    /// The current leader of the partition.
254    pub current_leader: LeaderIdAndEpoch,
255    /// The record data.
256    pub records: Option<Bytes>,
257    /// The acquired records.
258    pub acquired_records: Vec<AcquiredRecords>,
259    pub _unknown_tagged_fields: Vec<RawTaggedField>,
260}
261impl Default for PartitionData {
262    fn default() -> Self {
263        Self {
264            partition_index: 0_i32,
265            error_code: 0_i16,
266            error_message: None,
267            acknowledge_error_code: 0_i16,
268            acknowledge_error_message: None,
269            current_leader: LeaderIdAndEpoch::default(),
270            records: None,
271            acquired_records: Vec::new(),
272            _unknown_tagged_fields: Vec::new(),
273        }
274    }
275}
276impl PartitionData {
277    pub fn with_partition_index(mut self, value: i32) -> Self {
278        self.partition_index = value;
279        self
280    }
281    pub fn with_error_code(mut self, value: i16) -> Self {
282        self.error_code = value;
283        self
284    }
285    pub fn with_error_message(mut self, value: Option<KafkaString>) -> Self {
286        self.error_message = value;
287        self
288    }
289    pub fn with_acknowledge_error_code(mut self, value: i16) -> Self {
290        self.acknowledge_error_code = value;
291        self
292    }
293    pub fn with_acknowledge_error_message(mut self, value: Option<KafkaString>) -> Self {
294        self.acknowledge_error_message = value;
295        self
296    }
297    pub fn with_current_leader(mut self, value: LeaderIdAndEpoch) -> Self {
298        self.current_leader = value;
299        self
300    }
301    pub fn with_records(mut self, value: Option<Bytes>) -> Self {
302        self.records = value;
303        self
304    }
305    pub fn with_acquired_records(mut self, value: Vec<AcquiredRecords>) -> Self {
306        self.acquired_records = value;
307        self
308    }
309    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
310        let partition_index;
311        let error_code;
312        let error_message;
313        let acknowledge_error_code;
314        let acknowledge_error_message;
315        let current_leader;
316        let records;
317        let acquired_records;
318        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
319        partition_index = read_i32(buf)?;
320        error_code = read_i16(buf)?;
321        error_message = read_compact_nullable_string(buf)?;
322        acknowledge_error_code = read_i16(buf)?;
323        acknowledge_error_message = read_compact_nullable_string(buf)?;
324        current_leader = LeaderIdAndEpoch::read(buf, version)?;
325        records = read_compact_nullable_bytes(buf)?;
326        acquired_records = {
327            let len = read_compact_array_length(buf)?;
328            let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
329            for _ in 0..len {
330                arr.push(AcquiredRecords::read(buf, version)?);
331            }
332            arr
333        };
334        let tagged_fields = read_tagged_fields(buf)?;
335        for field in &tagged_fields {
336            match field.tag {
337                _ => {
338                    _unknown_tagged_fields.push(field.clone());
339                },
340            }
341        }
342        Ok(Self {
343            partition_index,
344            error_code,
345            error_message,
346            acknowledge_error_code,
347            acknowledge_error_message,
348            current_leader,
349            records,
350            acquired_records,
351            _unknown_tagged_fields,
352        })
353    }
354    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
355        write_i32(buf, self.partition_index);
356        write_i16(buf, self.error_code);
357        write_compact_nullable_string(buf, self.error_message.as_ref())?;
358        write_i16(buf, self.acknowledge_error_code);
359        write_compact_nullable_string(buf, self.acknowledge_error_message.as_ref())?;
360        self.current_leader.write(buf, version)?;
361        write_compact_nullable_bytes(buf, self.records.as_ref().map(|b| b.as_ref()))?;
362        write_compact_array_length(buf, self.acquired_records.len() as i32);
363        for el in &self.acquired_records {
364            el.write(buf, version)?;
365        }
366        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
367        all_tags.sort_by_key(|f| f.tag);
368        write_tagged_fields(buf, &all_tags)?;
369        Ok(())
370    }
371    pub fn encoded_len(&self, version: i16) -> Result<usize> {
372        let mut len: usize = 0;
373        len += 4;
374        len += 2;
375        len += compact_nullable_string_len(self.error_message.as_ref())?;
376        len += 2;
377        len += compact_nullable_string_len(self.acknowledge_error_message.as_ref())?;
378        len += self.current_leader.encoded_len(version)?;
379        len += compact_nullable_bytes_len(self.records.as_ref().map(|b| b.as_ref()))?;
380        len += compact_array_length_len(self.acquired_records.len() as i32);
381        for el in &self.acquired_records {
382            len += el.encoded_len(version)?;
383        }
384        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
385        all_tags.sort_by_key(|f| f.tag);
386        len += tagged_fields_len(&all_tags)?;
387        Ok(len)
388    }
389}
390#[derive(Debug, Clone, PartialEq)]
391pub struct LeaderIdAndEpoch {
392    /// The ID of the current leader or -1 if the leader is unknown.
393    pub leader_id: i32,
394    /// The latest known leader epoch.
395    pub leader_epoch: i32,
396    pub _unknown_tagged_fields: Vec<RawTaggedField>,
397}
398impl Default for LeaderIdAndEpoch {
399    fn default() -> Self {
400        Self {
401            leader_id: 0_i32,
402            leader_epoch: 0_i32,
403            _unknown_tagged_fields: Vec::new(),
404        }
405    }
406}
407impl LeaderIdAndEpoch {
408    pub fn with_leader_id(mut self, value: i32) -> Self {
409        self.leader_id = value;
410        self
411    }
412    pub fn with_leader_epoch(mut self, value: i32) -> Self {
413        self.leader_epoch = value;
414        self
415    }
416    pub fn read(buf: &mut Bytes, _version: i16) -> Result<Self> {
417        let leader_id;
418        let leader_epoch;
419        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
420        leader_id = read_i32(buf)?;
421        leader_epoch = read_i32(buf)?;
422        let tagged_fields = read_tagged_fields(buf)?;
423        for field in &tagged_fields {
424            match field.tag {
425                _ => {
426                    _unknown_tagged_fields.push(field.clone());
427                },
428            }
429        }
430        Ok(Self {
431            leader_id,
432            leader_epoch,
433            _unknown_tagged_fields,
434        })
435    }
436    pub fn write(&self, buf: &mut BytesMut, _version: i16) -> Result<()> {
437        write_i32(buf, self.leader_id);
438        write_i32(buf, self.leader_epoch);
439        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
440        all_tags.sort_by_key(|f| f.tag);
441        write_tagged_fields(buf, &all_tags)?;
442        Ok(())
443    }
444    pub fn encoded_len(&self, _version: i16) -> Result<usize> {
445        let mut len: usize = 0;
446        len += 4;
447        len += 4;
448        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
449        all_tags.sort_by_key(|f| f.tag);
450        len += tagged_fields_len(&all_tags)?;
451        Ok(len)
452    }
453}
454#[derive(Debug, Clone, PartialEq)]
455pub struct AcquiredRecords {
456    /// The earliest offset in this batch of acquired records.
457    pub first_offset: i64,
458    /// The last offset of this batch of acquired records.
459    pub last_offset: i64,
460    /// The delivery count of this batch of acquired records.
461    pub delivery_count: i16,
462    pub _unknown_tagged_fields: Vec<RawTaggedField>,
463}
464impl Default for AcquiredRecords {
465    fn default() -> Self {
466        Self {
467            first_offset: 0_i64,
468            last_offset: 0_i64,
469            delivery_count: 0_i16,
470            _unknown_tagged_fields: Vec::new(),
471        }
472    }
473}
474impl AcquiredRecords {
475    pub fn with_first_offset(mut self, value: i64) -> Self {
476        self.first_offset = value;
477        self
478    }
479    pub fn with_last_offset(mut self, value: i64) -> Self {
480        self.last_offset = value;
481        self
482    }
483    pub fn with_delivery_count(mut self, value: i16) -> Self {
484        self.delivery_count = value;
485        self
486    }
487    pub fn read(buf: &mut Bytes, _version: i16) -> Result<Self> {
488        let first_offset;
489        let last_offset;
490        let delivery_count;
491        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
492        first_offset = read_i64(buf)?;
493        last_offset = read_i64(buf)?;
494        delivery_count = read_i16(buf)?;
495        let tagged_fields = read_tagged_fields(buf)?;
496        for field in &tagged_fields {
497            match field.tag {
498                _ => {
499                    _unknown_tagged_fields.push(field.clone());
500                },
501            }
502        }
503        Ok(Self {
504            first_offset,
505            last_offset,
506            delivery_count,
507            _unknown_tagged_fields,
508        })
509    }
510    pub fn write(&self, buf: &mut BytesMut, _version: i16) -> Result<()> {
511        write_i64(buf, self.first_offset);
512        write_i64(buf, self.last_offset);
513        write_i16(buf, self.delivery_count);
514        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
515        all_tags.sort_by_key(|f| f.tag);
516        write_tagged_fields(buf, &all_tags)?;
517        Ok(())
518    }
519    pub fn encoded_len(&self, _version: i16) -> Result<usize> {
520        let mut len: usize = 0;
521        len += 8;
522        len += 8;
523        len += 2;
524        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
525        all_tags.sort_by_key(|f| f.tag);
526        len += tagged_fields_len(&all_tags)?;
527        Ok(len)
528    }
529}
530#[derive(Debug, Clone, PartialEq)]
531pub struct NodeEndpoint {
532    /// The ID of the associated node.
533    pub node_id: i32,
534    /// The node's hostname.
535    pub host: KafkaString,
536    /// The node's port.
537    pub port: i32,
538    /// The rack of the node, or null if it has not been assigned to a rack.
539    pub rack: Option<KafkaString>,
540    pub _unknown_tagged_fields: Vec<RawTaggedField>,
541}
542impl Default for NodeEndpoint {
543    fn default() -> Self {
544        Self {
545            node_id: 0_i32,
546            host: KafkaString::default(),
547            port: 0_i32,
548            rack: None,
549            _unknown_tagged_fields: Vec::new(),
550        }
551    }
552}
553impl NodeEndpoint {
554    pub fn with_node_id(mut self, value: i32) -> Self {
555        self.node_id = value;
556        self
557    }
558    pub fn with_host(mut self, value: KafkaString) -> Self {
559        self.host = value;
560        self
561    }
562    pub fn with_port(mut self, value: i32) -> Self {
563        self.port = value;
564        self
565    }
566    pub fn with_rack(mut self, value: Option<KafkaString>) -> Self {
567        self.rack = value;
568        self
569    }
570    pub fn read(buf: &mut Bytes, _version: i16) -> Result<Self> {
571        let node_id;
572        let host;
573        let port;
574        let rack;
575        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
576        node_id = read_i32(buf)?;
577        host = read_compact_string(buf)?;
578        port = read_i32(buf)?;
579        rack = read_compact_nullable_string(buf)?;
580        let tagged_fields = read_tagged_fields(buf)?;
581        for field in &tagged_fields {
582            match field.tag {
583                _ => {
584                    _unknown_tagged_fields.push(field.clone());
585                },
586            }
587        }
588        Ok(Self {
589            node_id,
590            host,
591            port,
592            rack,
593            _unknown_tagged_fields,
594        })
595    }
596    pub fn write(&self, buf: &mut BytesMut, _version: i16) -> Result<()> {
597        write_i32(buf, self.node_id);
598        write_compact_string(buf, &self.host)?;
599        write_i32(buf, self.port);
600        write_compact_nullable_string(buf, self.rack.as_ref())?;
601        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
602        all_tags.sort_by_key(|f| f.tag);
603        write_tagged_fields(buf, &all_tags)?;
604        Ok(())
605    }
606    pub fn encoded_len(&self, _version: i16) -> Result<usize> {
607        let mut len: usize = 0;
608        len += 4;
609        len += compact_string_len(&self.host)?;
610        len += 4;
611        len += compact_nullable_string_len(self.rack.as_ref())?;
612        let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
613        all_tags.sort_by_key(|f| f.tag);
614        len += tagged_fields_len(&all_tags)?;
615        Ok(len)
616    }
617}