Skip to main content

kacrab_protocol/generated/
list_offsets_request.rs

1//! Generated from ListOffsetsRequest.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 ListOffsetsRequestData {
17    /// The broker ID of the requester, or -1 if this request is being made by a normal consumer.
18    pub replica_id: i32,
19    /// This setting controls the visibility of transactional records. Using READ_UNCOMMITTED
20    /// (isolation_level = 0) makes all records visible. With READ_COMMITTED (isolation_level = 1),
21    /// non-transactional and COMMITTED transactional records are visible. To be more concrete,
22    /// READ_COMMITTED returns all data from offsets smaller than the current LSO (last stable
23    /// offset), and enables the inclusion of the list of aborted transactions in the result, which
24    /// allows consumers to discard ABORTED transactional records.
25    pub isolation_level: i8,
26    /// Each topic in the request.
27    pub topics: Vec<ListOffsetsTopic>,
28    /// The timeout to await a response in milliseconds for requests that require reading from
29    /// remote storage for topics enabled with tiered storage.
30    pub timeout_ms: i32,
31    pub _unknown_tagged_fields: Vec<RawTaggedField>,
32}
33impl Default for ListOffsetsRequestData {
34    fn default() -> Self {
35        Self {
36            replica_id: 0_i32,
37            isolation_level: 0_i8,
38            topics: Vec::new(),
39            timeout_ms: 0_i32,
40            _unknown_tagged_fields: Vec::new(),
41        }
42    }
43}
44impl ListOffsetsRequestData {
45    pub fn with_replica_id(mut self, value: i32) -> Self {
46        self.replica_id = value;
47        self
48    }
49    pub fn with_isolation_level(mut self, value: i8) -> Self {
50        self.isolation_level = value;
51        self
52    }
53    pub fn with_topics(mut self, value: Vec<ListOffsetsTopic>) -> Self {
54        self.topics = value;
55        self
56    }
57    pub fn with_timeout_ms(mut self, value: i32) -> Self {
58        self.timeout_ms = value;
59        self
60    }
61    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
62        if version < 1 || version > 11 {
63            return Err(UnsupportedVersion::new(2, version).into());
64        }
65        let replica_id;
66        let mut isolation_level = 0_i8;
67        let topics;
68        let mut timeout_ms = 0_i32;
69        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
70        replica_id = read_i32(buf)?;
71        if version >= 2 {
72            isolation_level = read_i8(buf)?;
73        }
74        if version >= 6 {
75            topics = {
76                let len = read_compact_array_length(buf)?;
77                let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
78                for _ in 0..len {
79                    arr.push(ListOffsetsTopic::read(buf, version)?);
80                }
81                arr
82            };
83        } else {
84            topics = {
85                let len = read_array_length(buf)?;
86                let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
87                for _ in 0..len {
88                    arr.push(ListOffsetsTopic::read(buf, version)?);
89                }
90                arr
91            };
92        }
93        if version >= 10 {
94            timeout_ms = read_i32(buf)?;
95        }
96        if version >= 6 {
97            let tagged_fields = read_tagged_fields(buf)?;
98            for field in &tagged_fields {
99                match field.tag {
100                    _ => {
101                        _unknown_tagged_fields.push(field.clone());
102                    },
103                }
104            }
105        }
106        Ok(Self {
107            replica_id,
108            isolation_level,
109            topics,
110            timeout_ms,
111            _unknown_tagged_fields,
112        })
113    }
114    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
115        if version < 1 || version > 11 {
116            return Err(UnsupportedVersion::new(2, version).into());
117        }
118        write_i32(buf, self.replica_id);
119        if version >= 2 {
120            write_i8(buf, self.isolation_level);
121        } else if self.isolation_level != 0_i8 {
122            return Err(UnsupportedFieldVersion::new(2, "isolation_level", version).into());
123        }
124        if version >= 6 {
125            write_compact_array_length(buf, self.topics.len() as i32);
126            for el in &self.topics {
127                el.write(buf, version)?;
128            }
129        } else {
130            write_array_length(buf, self.topics.len() as i32);
131            for el in &self.topics {
132                el.write(buf, version)?;
133            }
134        }
135        if version >= 10 {
136            write_i32(buf, self.timeout_ms);
137        } else if self.timeout_ms != 0_i32 {
138            return Err(UnsupportedFieldVersion::new(2, "timeout_ms", version).into());
139        }
140        if version >= 6 {
141            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
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 < 1 || version > 11 {
149            return Err(UnsupportedVersion::new(2, version).into());
150        }
151        let mut len: usize = 0;
152        len += 4;
153        if version >= 2 {
154            len += 1;
155        } else if self.isolation_level != 0_i8 {
156            return Err(UnsupportedFieldVersion::new(2, "isolation_level", version).into());
157        }
158        if version >= 6 {
159            len += compact_array_length_len(self.topics.len() as i32);
160            for el in &self.topics {
161                len += el.encoded_len(version)?;
162            }
163        } else {
164            len += array_length_len();
165            for el in &self.topics {
166                len += el.encoded_len(version)?;
167            }
168        }
169        if version >= 10 {
170            len += 4;
171        } else if self.timeout_ms != 0_i32 {
172            return Err(UnsupportedFieldVersion::new(2, "timeout_ms", version).into());
173        }
174        if version >= 6 {
175            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
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 ListOffsetsTopic {
184    /// The topic name.
185    pub name: KafkaString,
186    /// Each partition in the request.
187    pub partitions: Vec<ListOffsetsPartition>,
188    pub _unknown_tagged_fields: Vec<RawTaggedField>,
189}
190impl Default for ListOffsetsTopic {
191    fn default() -> Self {
192        Self {
193            name: KafkaString::default(),
194            partitions: Vec::new(),
195            _unknown_tagged_fields: Vec::new(),
196        }
197    }
198}
199impl ListOffsetsTopic {
200    pub fn with_name(mut self, value: KafkaString) -> Self {
201        self.name = value;
202        self
203    }
204    pub fn with_partitions(mut self, value: Vec<ListOffsetsPartition>) -> Self {
205        self.partitions = value;
206        self
207    }
208    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
209        let name;
210        let partitions;
211        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
212        if version >= 6 {
213            name = read_compact_string(buf)?;
214        } else {
215            name = read_string(buf)?;
216        }
217        if version >= 6 {
218            partitions = {
219                let len = read_compact_array_length(buf)?;
220                let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
221                for _ in 0..len {
222                    arr.push(ListOffsetsPartition::read(buf, version)?);
223                }
224                arr
225            };
226        } else {
227            partitions = {
228                let len = read_array_length(buf)?;
229                let mut arr = Vec::with_capacity(array_read_capacity(len, (buf).len()));
230                for _ in 0..len {
231                    arr.push(ListOffsetsPartition::read(buf, version)?);
232                }
233                arr
234            };
235        }
236        if version >= 6 {
237            let tagged_fields = read_tagged_fields(buf)?;
238            for field in &tagged_fields {
239                match field.tag {
240                    _ => {
241                        _unknown_tagged_fields.push(field.clone());
242                    },
243                }
244            }
245        }
246        Ok(Self {
247            name,
248            partitions,
249            _unknown_tagged_fields,
250        })
251    }
252    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
253        if version >= 6 {
254            write_compact_string(buf, &self.name)?;
255        } else {
256            write_string(buf, &self.name)?;
257        }
258        if version >= 6 {
259            write_compact_array_length(buf, self.partitions.len() as i32);
260            for el in &self.partitions {
261                el.write(buf, version)?;
262            }
263        } else {
264            write_array_length(buf, self.partitions.len() as i32);
265            for el in &self.partitions {
266                el.write(buf, version)?;
267            }
268        }
269        if version >= 6 {
270            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
271            all_tags.sort_by_key(|f| f.tag);
272            write_tagged_fields(buf, &all_tags)?;
273        }
274        Ok(())
275    }
276    pub fn encoded_len(&self, version: i16) -> Result<usize> {
277        let mut len: usize = 0;
278        if version >= 6 {
279            len += compact_string_len(&self.name)?;
280        } else {
281            len += string_len(&self.name)?;
282        }
283        if version >= 6 {
284            len += compact_array_length_len(self.partitions.len() as i32);
285            for el in &self.partitions {
286                len += el.encoded_len(version)?;
287            }
288        } else {
289            len += array_length_len();
290            for el in &self.partitions {
291                len += el.encoded_len(version)?;
292            }
293        }
294        if version >= 6 {
295            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
296            all_tags.sort_by_key(|f| f.tag);
297            len += tagged_fields_len(&all_tags)?;
298        }
299        Ok(len)
300    }
301}
302#[derive(Debug, Clone, PartialEq)]
303pub struct ListOffsetsPartition {
304    /// The partition index.
305    pub partition_index: i32,
306    /// The current leader epoch.
307    pub current_leader_epoch: i32,
308    /// The current timestamp.
309    pub timestamp: i64,
310    pub _unknown_tagged_fields: Vec<RawTaggedField>,
311}
312impl Default for ListOffsetsPartition {
313    fn default() -> Self {
314        Self {
315            partition_index: 0_i32,
316            current_leader_epoch: -1i32,
317            timestamp: 0_i64,
318            _unknown_tagged_fields: Vec::new(),
319        }
320    }
321}
322impl ListOffsetsPartition {
323    pub fn with_partition_index(mut self, value: i32) -> Self {
324        self.partition_index = value;
325        self
326    }
327    pub fn with_current_leader_epoch(mut self, value: i32) -> Self {
328        self.current_leader_epoch = value;
329        self
330    }
331    pub fn with_timestamp(mut self, value: i64) -> Self {
332        self.timestamp = value;
333        self
334    }
335    pub fn read(buf: &mut Bytes, version: i16) -> Result<Self> {
336        let partition_index;
337        let mut current_leader_epoch = -1i32;
338        let timestamp;
339        let mut _unknown_tagged_fields: Vec<RawTaggedField> = Vec::new();
340        partition_index = read_i32(buf)?;
341        if version >= 4 {
342            current_leader_epoch = read_i32(buf)?;
343        }
344        timestamp = read_i64(buf)?;
345        if version >= 6 {
346            let tagged_fields = read_tagged_fields(buf)?;
347            for field in &tagged_fields {
348                match field.tag {
349                    _ => {
350                        _unknown_tagged_fields.push(field.clone());
351                    },
352                }
353            }
354        }
355        Ok(Self {
356            partition_index,
357            current_leader_epoch,
358            timestamp,
359            _unknown_tagged_fields,
360        })
361    }
362    pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
363        write_i32(buf, self.partition_index);
364        if version >= 4 {
365            write_i32(buf, self.current_leader_epoch);
366        } else if self.current_leader_epoch != -1i32 {
367            return Err(UnsupportedFieldVersion::new(2, "current_leader_epoch", version).into());
368        }
369        write_i64(buf, self.timestamp);
370        if version >= 6 {
371            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
372            all_tags.sort_by_key(|f| f.tag);
373            write_tagged_fields(buf, &all_tags)?;
374        }
375        Ok(())
376    }
377    pub fn encoded_len(&self, version: i16) -> Result<usize> {
378        let mut len: usize = 0;
379        len += 4;
380        if version >= 4 {
381            len += 4;
382        } else if self.current_leader_epoch != -1i32 {
383            return Err(UnsupportedFieldVersion::new(2, "current_leader_epoch", version).into());
384        }
385        len += 8;
386        if version >= 6 {
387            let mut all_tags: Vec<RawTaggedField> = self._unknown_tagged_fields.clone();
388            all_tags.sort_by_key(|f| f.tag);
389            len += tagged_fields_len(&all_tags)?;
390        }
391        Ok(len)
392    }
393}