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 ShareFetchResponseData {
17 pub throttle_time_ms: i32,
20 pub error_code: i16,
22 pub error_message: Option<KafkaString>,
24 pub acquisition_lock_timeout_ms: i32,
26 pub responses: Vec<ShareFetchableTopicResponse>,
28 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 pub topic_id: KafkaUuid,
168 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 pub partition_index: i32,
245 pub error_code: i16,
247 pub error_message: Option<KafkaString>,
249 pub acknowledge_error_code: i16,
251 pub acknowledge_error_message: Option<KafkaString>,
253 pub current_leader: LeaderIdAndEpoch,
255 pub records: Option<Bytes>,
257 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 pub leader_id: i32,
394 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 pub first_offset: i64,
458 pub last_offset: i64,
460 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 pub node_id: i32,
534 pub host: KafkaString,
536 pub port: i32,
538 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}