1use crate::codec::{Decoder, Encoder};
2use crate::error::{Error, Result};
3use crate::header::RequestHeader;
4
5pub const SHARE_GROUP_HEARTBEAT_API_KEY: i16 = 76;
7pub const SHARE_FETCH_API_KEY: i16 = 78;
9pub const SHARE_ACKNOWLEDGE_API_KEY: i16 = 79;
11
12#[derive(Debug, Clone, PartialEq, Eq)]
14pub struct ShareTopicPartitionsV1 {
15 pub topic_id: [u8; 16],
16 pub partitions: Vec<i32>,
17}
18
19impl ShareTopicPartitionsV1 {
20 #[cfg(test)]
21 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
22 encoder.write_uuid(&self.topic_id);
23 encoder.write_compact_array(Some(&self.partitions), |encoder, partition| {
24 encoder.write_i32(*partition);
25 Ok(())
26 })?;
27 encoder.write_empty_tagged_fields();
28 Ok(())
29 }
30
31 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
32 let topic_id = decoder.read_uuid()?;
33 let partitions = decoder
34 .read_compact_array("share topic partitions", |decoder| decoder.read_i32())?
35 .unwrap_or_default();
36 decoder.read_tagged_fields()?;
37 Ok(Self {
38 topic_id,
39 partitions,
40 })
41 }
42}
43
44#[derive(Debug, Clone, PartialEq, Eq)]
49pub struct ShareGroupHeartbeatRequestV1 {
50 pub correlation_id: i32,
51 pub client_id: Option<String>,
52 pub group_id: String,
53 pub member_id: String,
54 pub member_epoch: i32,
55 pub rack_id: Option<String>,
56 pub subscribed_topic_names: Option<Vec<String>>,
57}
58
59impl ShareGroupHeartbeatRequestV1 {
60 pub fn encode(&self) -> Result<Vec<u8>> {
62 let mut encoder = Encoder::new();
63 RequestHeader {
64 api_key: SHARE_GROUP_HEARTBEAT_API_KEY,
65 api_version: 1,
66 correlation_id: self.correlation_id,
67 client_id: self.client_id.clone(),
68 }
69 .encode_v2(&mut encoder)?;
70 encoder.write_compact_string(&self.group_id)?;
71 encoder.write_compact_string(&self.member_id)?;
72 encoder.write_i32(self.member_epoch);
73 encoder.write_compact_nullable_string(self.rack_id.as_deref())?;
74 encoder.write_compact_array(self.subscribed_topic_names.as_deref(), |encoder, topic| {
75 encoder.write_compact_string(topic)
76 })?;
77 encoder.write_empty_tagged_fields();
78 Ok(encoder.into_bytes())
79 }
80}
81
82#[derive(Debug, Clone, PartialEq, Eq)]
84pub struct ShareGroupHeartbeatAssignmentV1 {
85 pub topic_partitions: Vec<ShareTopicPartitionsV1>,
86}
87
88impl ShareGroupHeartbeatAssignmentV1 {
89 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
90 let topic_partitions = decoder
91 .read_compact_array("share heartbeat assignment", ShareTopicPartitionsV1::decode)?
92 .unwrap_or_default();
93 decoder.read_tagged_fields()?;
94 Ok(Self { topic_partitions })
95 }
96}
97
98#[derive(Debug, Clone, PartialEq, Eq)]
100pub struct ShareGroupHeartbeatResponseV1 {
101 pub throttle_time_ms: i32,
102 pub error_code: i16,
103 pub error_message: Option<String>,
104 pub member_id: Option<String>,
105 pub member_epoch: i32,
106 pub heartbeat_interval_ms: i32,
107 pub assignment: Option<ShareGroupHeartbeatAssignmentV1>,
108}
109
110impl ShareGroupHeartbeatResponseV1 {
111 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
113 let throttle_time_ms = decoder.read_i32()?;
114 let error_code = decoder.read_i16()?;
115 let error_message = decoder.read_compact_nullable_string()?;
116 let member_id = decoder.read_compact_nullable_string()?;
117 let member_epoch = decoder.read_i32()?;
118 let heartbeat_interval_ms = decoder.read_i32()?;
119 let assignment = match decoder.read_i8()? {
120 -1 => None,
121 1 => Some(ShareGroupHeartbeatAssignmentV1::decode(decoder)?),
122 marker => return Err(Error::InvalidNullableStruct(marker)),
123 };
124 decoder.read_tagged_fields()?;
125 Ok(Self {
126 throttle_time_ms,
127 error_code,
128 error_message,
129 member_id,
130 member_epoch,
131 heartbeat_interval_ms,
132 assignment,
133 })
134 }
135}
136
137#[derive(Debug, Clone, PartialEq, Eq)]
139pub struct ShareAcknowledgementBatchV1 {
140 pub first_offset: i64,
141 pub last_offset: i64,
142 pub acknowledgement_types: Vec<i8>,
144}
145
146impl ShareAcknowledgementBatchV1 {
147 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
148 encoder.write_i64(self.first_offset);
149 encoder.write_i64(self.last_offset);
150 encoder.write_compact_array(
151 Some(&self.acknowledgement_types),
152 |encoder, acknowledgement_type| {
153 encoder.write_i8(*acknowledgement_type);
154 Ok(())
155 },
156 )?;
157 encoder.write_empty_tagged_fields();
158 Ok(())
159 }
160
161 #[cfg(test)]
162 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
163 let first_offset = decoder.read_i64()?;
164 let last_offset = decoder.read_i64()?;
165 let acknowledgement_types = decoder
166 .read_compact_array("share acknowledgement types", |decoder| decoder.read_i8())?
167 .unwrap_or_default();
168 decoder.read_tagged_fields()?;
169 Ok(Self {
170 first_offset,
171 last_offset,
172 acknowledgement_types,
173 })
174 }
175}
176
177#[derive(Debug, Clone, PartialEq, Eq)]
179pub struct ShareFetchPartitionV1 {
180 pub partition_index: i32,
181 pub acknowledgement_batches: Vec<ShareAcknowledgementBatchV1>,
182}
183
184impl ShareFetchPartitionV1 {
185 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
186 encoder.write_i32(self.partition_index);
187 encoder.write_compact_array(Some(&self.acknowledgement_batches), |encoder, batch| {
188 batch.encode(encoder)
189 })?;
190 encoder.write_empty_tagged_fields();
191 Ok(())
192 }
193}
194
195#[derive(Debug, Clone, PartialEq, Eq)]
197pub struct ShareFetchTopicV1 {
198 pub topic_id: [u8; 16],
199 pub partitions: Vec<ShareFetchPartitionV1>,
200}
201
202impl ShareFetchTopicV1 {
203 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
204 encoder.write_uuid(&self.topic_id);
205 encoder.write_compact_array(Some(&self.partitions), |encoder, partition| {
206 partition.encode(encoder)
207 })?;
208 encoder.write_empty_tagged_fields();
209 Ok(())
210 }
211}
212
213#[derive(Debug, Clone, PartialEq, Eq)]
215pub struct ShareForgottenTopicV1 {
216 pub topic_id: [u8; 16],
217 pub partitions: Vec<i32>,
218}
219
220impl ShareForgottenTopicV1 {
221 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
222 encoder.write_uuid(&self.topic_id);
223 encoder.write_compact_array(Some(&self.partitions), |encoder, partition| {
224 encoder.write_i32(*partition);
225 Ok(())
226 })?;
227 encoder.write_empty_tagged_fields();
228 Ok(())
229 }
230}
231
232#[derive(Debug, Clone, PartialEq, Eq)]
234pub struct ShareFetchRequestV1 {
235 pub correlation_id: i32,
236 pub client_id: Option<String>,
237 pub group_id: Option<String>,
238 pub member_id: Option<String>,
239 pub share_session_epoch: i32,
240 pub max_wait_ms: i32,
241 pub min_bytes: i32,
242 pub max_bytes: i32,
243 pub max_records: i32,
244 pub batch_size: i32,
245 pub topics: Vec<ShareFetchTopicV1>,
246 pub forgotten_topics: Vec<ShareForgottenTopicV1>,
247}
248
249impl ShareFetchRequestV1 {
250 pub fn encode(&self) -> Result<Vec<u8>> {
252 encode_share_fetch_request(
253 1,
254 self.correlation_id,
255 self.client_id.clone(),
256 self.group_id.clone(),
257 self.member_id.clone(),
258 self.share_session_epoch,
259 self.max_wait_ms,
260 self.min_bytes,
261 self.max_bytes,
262 self.max_records,
263 self.batch_size,
264 None,
265 false,
266 &self.topics,
267 &self.forgotten_topics,
268 )
269 }
270}
271
272#[derive(Debug, Clone, PartialEq, Eq)]
276pub struct ShareFetchRequestV2 {
277 pub correlation_id: i32,
278 pub client_id: Option<String>,
279 pub group_id: Option<String>,
280 pub member_id: Option<String>,
281 pub share_session_epoch: i32,
282 pub max_wait_ms: i32,
283 pub min_bytes: i32,
284 pub max_bytes: i32,
285 pub max_records: i32,
286 pub batch_size: i32,
287 pub share_acquire_mode: i8,
289 pub is_renew_ack: bool,
291 pub topics: Vec<ShareFetchTopicV1>,
292 pub forgotten_topics: Vec<ShareForgottenTopicV1>,
293}
294
295impl ShareFetchRequestV2 {
296 pub fn encode(&self) -> Result<Vec<u8>> {
298 encode_share_fetch_request(
299 2,
300 self.correlation_id,
301 self.client_id.clone(),
302 self.group_id.clone(),
303 self.member_id.clone(),
304 self.share_session_epoch,
305 self.max_wait_ms,
306 self.min_bytes,
307 self.max_bytes,
308 self.max_records,
309 self.batch_size,
310 Some(self.share_acquire_mode),
311 self.is_renew_ack,
312 &self.topics,
313 &self.forgotten_topics,
314 )
315 }
316}
317
318#[allow(clippy::too_many_arguments)]
319fn encode_share_fetch_request(
320 api_version: i16,
321 correlation_id: i32,
322 client_id: Option<String>,
323 group_id: Option<String>,
324 member_id: Option<String>,
325 share_session_epoch: i32,
326 max_wait_ms: i32,
327 min_bytes: i32,
328 max_bytes: i32,
329 max_records: i32,
330 batch_size: i32,
331 share_acquire_mode: Option<i8>,
332 is_renew_ack: bool,
333 topics: &[ShareFetchTopicV1],
334 forgotten_topics: &[ShareForgottenTopicV1],
335) -> Result<Vec<u8>> {
336 let mut encoder = Encoder::new();
337 RequestHeader {
338 api_key: SHARE_FETCH_API_KEY,
339 api_version,
340 correlation_id,
341 client_id,
342 }
343 .encode_v2(&mut encoder)?;
344 encoder.write_compact_nullable_string(group_id.as_deref())?;
345 encoder.write_compact_nullable_string(member_id.as_deref())?;
346 encoder.write_i32(share_session_epoch);
347 encoder.write_i32(max_wait_ms);
348 encoder.write_i32(min_bytes);
349 encoder.write_i32(max_bytes);
350 encoder.write_i32(max_records);
351 encoder.write_i32(batch_size);
352 if let Some(share_acquire_mode) = share_acquire_mode {
353 encoder.write_i8(share_acquire_mode);
354 encoder.write_bool(is_renew_ack);
355 }
356 encoder.write_compact_array(Some(topics), |encoder, topic| topic.encode(encoder))?;
357 encoder.write_compact_array(Some(forgotten_topics), |encoder, topic| {
358 topic.encode(encoder)
359 })?;
360 encoder.write_empty_tagged_fields();
361 Ok(encoder.into_bytes())
362}
363
364#[derive(Debug, Clone, PartialEq, Eq)]
366pub struct ShareLeaderIdAndEpochV1 {
367 pub leader_id: i32,
368 pub leader_epoch: i32,
369}
370
371impl ShareLeaderIdAndEpochV1 {
372 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
373 let leader_id = decoder.read_i32()?;
374 let leader_epoch = decoder.read_i32()?;
375 decoder.read_tagged_fields()?;
376 Ok(Self {
377 leader_id,
378 leader_epoch,
379 })
380 }
381}
382
383#[derive(Debug, Clone, PartialEq, Eq)]
385pub struct ShareAcquiredRecordsV1 {
386 pub first_offset: i64,
387 pub last_offset: i64,
388 pub delivery_count: i16,
389}
390
391impl ShareAcquiredRecordsV1 {
392 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
393 let first_offset = decoder.read_i64()?;
394 let last_offset = decoder.read_i64()?;
395 let delivery_count = decoder.read_i16()?;
396 decoder.read_tagged_fields()?;
397 Ok(Self {
398 first_offset,
399 last_offset,
400 delivery_count,
401 })
402 }
403}
404
405#[derive(Debug, Clone, PartialEq, Eq)]
407pub struct ShareFetchPartitionResponseV1 {
408 pub partition_index: i32,
409 pub error_code: i16,
410 pub error_message: Option<String>,
411 pub acknowledgement_error_code: i16,
412 pub acknowledgement_error_message: Option<String>,
413 pub current_leader: ShareLeaderIdAndEpochV1,
414 pub records: Option<Vec<u8>>,
416 pub acquired_records: Vec<ShareAcquiredRecordsV1>,
417}
418
419impl ShareFetchPartitionResponseV1 {
420 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
421 let partition_index = decoder.read_i32()?;
422 let error_code = decoder.read_i16()?;
423 let error_message = decoder.read_compact_nullable_string()?;
424 let acknowledgement_error_code = decoder.read_i16()?;
425 let acknowledgement_error_message = decoder.read_compact_nullable_string()?;
426 let current_leader = ShareLeaderIdAndEpochV1::decode(decoder)?;
427 let records = decoder.read_compact_nullable_bytes()?;
428 let acquired_records = decoder
429 .read_compact_array("share acquired records", ShareAcquiredRecordsV1::decode)?
430 .unwrap_or_default();
431 decoder.read_tagged_fields()?;
432 Ok(Self {
433 partition_index,
434 error_code,
435 error_message,
436 acknowledgement_error_code,
437 acknowledgement_error_message,
438 current_leader,
439 records,
440 acquired_records,
441 })
442 }
443}
444
445#[derive(Debug, Clone, PartialEq, Eq)]
447pub struct ShareFetchTopicResponseV1 {
448 pub topic_id: [u8; 16],
449 pub partitions: Vec<ShareFetchPartitionResponseV1>,
450}
451
452impl ShareFetchTopicResponseV1 {
453 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
454 let topic_id = decoder.read_uuid()?;
455 let partitions = decoder
456 .read_compact_array(
457 "share fetch partitions",
458 ShareFetchPartitionResponseV1::decode,
459 )?
460 .unwrap_or_default();
461 decoder.read_tagged_fields()?;
462 Ok(Self {
463 topic_id,
464 partitions,
465 })
466 }
467}
468
469#[derive(Debug, Clone, PartialEq, Eq)]
471pub struct ShareNodeEndpointV1 {
472 pub node_id: i32,
473 pub host: String,
474 pub port: i32,
475 pub rack: Option<String>,
476}
477
478impl ShareNodeEndpointV1 {
479 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
480 let node_id = decoder.read_i32()?;
481 let host = decoder.read_compact_string()?;
482 let port = decoder.read_i32()?;
483 let rack = decoder.read_compact_nullable_string()?;
484 decoder.read_tagged_fields()?;
485 Ok(Self {
486 node_id,
487 host,
488 port,
489 rack,
490 })
491 }
492}
493
494#[derive(Debug, Clone, PartialEq, Eq)]
496pub struct ShareFetchResponseV1 {
497 pub throttle_time_ms: i32,
498 pub error_code: i16,
499 pub error_message: Option<String>,
500 pub acquisition_lock_timeout_ms: i32,
501 pub responses: Vec<ShareFetchTopicResponseV1>,
502 pub node_endpoints: Vec<ShareNodeEndpointV1>,
503}
504
505impl ShareFetchResponseV1 {
506 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
508 let throttle_time_ms = decoder.read_i32()?;
509 let error_code = decoder.read_i16()?;
510 let error_message = decoder.read_compact_nullable_string()?;
511 let acquisition_lock_timeout_ms = decoder.read_i32()?;
512 let responses = decoder
513 .read_compact_array("share fetch responses", ShareFetchTopicResponseV1::decode)?
514 .unwrap_or_default();
515 let node_endpoints = decoder
516 .read_compact_array("share node endpoints", ShareNodeEndpointV1::decode)?
517 .unwrap_or_default();
518 decoder.read_tagged_fields()?;
519 Ok(Self {
520 throttle_time_ms,
521 error_code,
522 error_message,
523 acquisition_lock_timeout_ms,
524 responses,
525 node_endpoints,
526 })
527 }
528}
529
530#[derive(Debug, Clone, PartialEq, Eq)]
532pub struct ShareAcknowledgeTopicV1 {
533 pub topic_id: [u8; 16],
534 pub partitions: Vec<ShareAcknowledgePartitionV1>,
535}
536
537impl ShareAcknowledgeTopicV1 {
538 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
539 encoder.write_uuid(&self.topic_id);
540 encoder.write_compact_array(Some(&self.partitions), |encoder, partition| {
541 partition.encode(encoder)
542 })?;
543 encoder.write_empty_tagged_fields();
544 Ok(())
545 }
546}
547
548#[derive(Debug, Clone, PartialEq, Eq)]
550pub struct ShareAcknowledgePartitionV1 {
551 pub partition_index: i32,
552 pub acknowledgement_batches: Vec<ShareAcknowledgementBatchV1>,
553}
554
555impl ShareAcknowledgePartitionV1 {
556 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
557 encoder.write_i32(self.partition_index);
558 encoder.write_compact_array(Some(&self.acknowledgement_batches), |encoder, batch| {
559 batch.encode(encoder)
560 })?;
561 encoder.write_empty_tagged_fields();
562 Ok(())
563 }
564}
565
566#[derive(Debug, Clone, PartialEq, Eq)]
568pub struct ShareAcknowledgeRequestV1 {
569 pub correlation_id: i32,
570 pub client_id: Option<String>,
571 pub group_id: Option<String>,
572 pub member_id: Option<String>,
573 pub share_session_epoch: i32,
574 pub topics: Vec<ShareAcknowledgeTopicV1>,
575}
576
577impl ShareAcknowledgeRequestV1 {
578 pub fn encode(&self) -> Result<Vec<u8>> {
580 encode_share_acknowledge_request(ShareAcknowledgeRequestParts {
581 api_version: 1,
582 correlation_id: self.correlation_id,
583 client_id: self.client_id.as_deref(),
584 group_id: self.group_id.as_deref(),
585 member_id: self.member_id.as_deref(),
586 share_session_epoch: self.share_session_epoch,
587 is_renew_ack: false,
588 topics: &self.topics,
589 })
590 }
591}
592
593#[derive(Debug, Clone, PartialEq, Eq)]
595pub struct ShareAcknowledgeRequestV2 {
596 pub correlation_id: i32,
597 pub client_id: Option<String>,
598 pub group_id: Option<String>,
599 pub member_id: Option<String>,
600 pub share_session_epoch: i32,
601 pub is_renew_ack: bool,
603 pub topics: Vec<ShareAcknowledgeTopicV1>,
604}
605
606impl ShareAcknowledgeRequestV2 {
607 pub fn encode(&self) -> Result<Vec<u8>> {
609 encode_share_acknowledge_request(ShareAcknowledgeRequestParts {
610 api_version: 2,
611 correlation_id: self.correlation_id,
612 client_id: self.client_id.as_deref(),
613 group_id: self.group_id.as_deref(),
614 member_id: self.member_id.as_deref(),
615 share_session_epoch: self.share_session_epoch,
616 is_renew_ack: self.is_renew_ack,
617 topics: &self.topics,
618 })
619 }
620}
621
622struct ShareAcknowledgeRequestParts<'a> {
623 api_version: i16,
624 correlation_id: i32,
625 client_id: Option<&'a str>,
626 group_id: Option<&'a str>,
627 member_id: Option<&'a str>,
628 share_session_epoch: i32,
629 is_renew_ack: bool,
630 topics: &'a [ShareAcknowledgeTopicV1],
631}
632
633fn encode_share_acknowledge_request(parts: ShareAcknowledgeRequestParts<'_>) -> Result<Vec<u8>> {
634 let mut encoder = Encoder::new();
635 RequestHeader {
636 api_key: SHARE_ACKNOWLEDGE_API_KEY,
637 api_version: parts.api_version,
638 correlation_id: parts.correlation_id,
639 client_id: parts.client_id.map(str::to_owned),
640 }
641 .encode_v2(&mut encoder)?;
642 encoder.write_compact_nullable_string(parts.group_id)?;
643 encoder.write_compact_nullable_string(parts.member_id)?;
644 encoder.write_i32(parts.share_session_epoch);
645 if parts.api_version >= 2 {
646 encoder.write_bool(parts.is_renew_ack);
647 }
648 encoder.write_compact_array(Some(parts.topics), |encoder, topic| topic.encode(encoder))?;
649 encoder.write_empty_tagged_fields();
650 Ok(encoder.into_bytes())
651}
652
653#[derive(Debug, Clone, PartialEq, Eq)]
655pub struct ShareAcknowledgePartitionResponseV1 {
656 pub partition_index: i32,
657 pub error_code: i16,
658 pub error_message: Option<String>,
659 pub current_leader: ShareLeaderIdAndEpochV1,
660}
661
662impl ShareAcknowledgePartitionResponseV1 {
663 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
664 let partition_index = decoder.read_i32()?;
665 let error_code = decoder.read_i16()?;
666 let error_message = decoder.read_compact_nullable_string()?;
667 let current_leader = ShareLeaderIdAndEpochV1::decode(decoder)?;
668 decoder.read_tagged_fields()?;
669 Ok(Self {
670 partition_index,
671 error_code,
672 error_message,
673 current_leader,
674 })
675 }
676}
677
678#[derive(Debug, Clone, PartialEq, Eq)]
680pub struct ShareAcknowledgeTopicResponseV1 {
681 pub topic_id: [u8; 16],
682 pub partitions: Vec<ShareAcknowledgePartitionResponseV1>,
683}
684
685impl ShareAcknowledgeTopicResponseV1 {
686 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
687 let topic_id = decoder.read_uuid()?;
688 let partitions = decoder
689 .read_compact_array(
690 "share acknowledgement partitions",
691 ShareAcknowledgePartitionResponseV1::decode,
692 )?
693 .unwrap_or_default();
694 decoder.read_tagged_fields()?;
695 Ok(Self {
696 topic_id,
697 partitions,
698 })
699 }
700}
701
702#[derive(Debug, Clone, PartialEq, Eq)]
704pub struct ShareAcknowledgeResponseV1 {
705 pub throttle_time_ms: i32,
706 pub error_code: i16,
707 pub error_message: Option<String>,
708 pub responses: Vec<ShareAcknowledgeTopicResponseV1>,
709 pub node_endpoints: Vec<ShareNodeEndpointV1>,
710}
711
712impl ShareAcknowledgeResponseV1 {
713 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
715 let throttle_time_ms = decoder.read_i32()?;
716 let error_code = decoder.read_i16()?;
717 let error_message = decoder.read_compact_nullable_string()?;
718 let responses = decoder
719 .read_compact_array(
720 "share acknowledgement responses",
721 ShareAcknowledgeTopicResponseV1::decode,
722 )?
723 .unwrap_or_default();
724 let node_endpoints = decoder
725 .read_compact_array("share node endpoints", ShareNodeEndpointV1::decode)?
726 .unwrap_or_default();
727 decoder.read_tagged_fields()?;
728 Ok(Self {
729 throttle_time_ms,
730 error_code,
731 error_message,
732 responses,
733 node_endpoints,
734 })
735 }
736}
737
738#[derive(Debug, Clone, PartialEq, Eq)]
740pub struct ShareAcknowledgeResponseV2 {
741 pub throttle_time_ms: i32,
742 pub error_code: i16,
743 pub error_message: Option<String>,
744 pub acquisition_lock_timeout_ms: i32,
745 pub responses: Vec<ShareAcknowledgeTopicResponseV1>,
746 pub node_endpoints: Vec<ShareNodeEndpointV1>,
747}
748
749impl ShareAcknowledgeResponseV2 {
750 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
752 let throttle_time_ms = decoder.read_i32()?;
753 let error_code = decoder.read_i16()?;
754 let error_message = decoder.read_compact_nullable_string()?;
755 let acquisition_lock_timeout_ms = decoder.read_i32()?;
756 let responses = decoder
757 .read_compact_array(
758 "share acknowledgement responses",
759 ShareAcknowledgeTopicResponseV1::decode,
760 )?
761 .unwrap_or_default();
762 let node_endpoints = decoder
763 .read_compact_array("share node endpoints", ShareNodeEndpointV1::decode)?
764 .unwrap_or_default();
765 decoder.read_tagged_fields()?;
766 Ok(Self {
767 throttle_time_ms,
768 error_code,
769 error_message,
770 acquisition_lock_timeout_ms,
771 responses,
772 node_endpoints,
773 })
774 }
775}
776
777#[cfg(test)]
778#[allow(clippy::unwrap_used)]
779mod tests {
780 use super::*;
781 use crate::codec::{Decoder, Encoder};
782
783 #[test]
784 fn encodes_share_group_heartbeat_v1_request() {
785 let request = ShareGroupHeartbeatRequestV1 {
786 correlation_id: 11,
787 client_id: Some("kafrust".to_owned()),
788 group_id: "share-orders".to_owned(),
789 member_id: "member-1".to_owned(),
790 member_epoch: 3,
791 rack_id: Some("rack-a".to_owned()),
792 subscribed_topic_names: Some(vec!["orders".to_owned()]),
793 };
794
795 let encoded = request.encode().unwrap();
796 let mut decoder = Decoder::new(&encoded);
797 assert_eq!(decoder.read_i16().unwrap(), SHARE_GROUP_HEARTBEAT_API_KEY);
798 assert_eq!(decoder.read_i16().unwrap(), 1);
799 assert_eq!(decoder.read_i32().unwrap(), 11);
800 assert_eq!(
801 decoder.read_nullable_string().unwrap(),
802 Some("kafrust".to_owned())
803 );
804 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
805 assert_eq!(decoder.read_compact_string().unwrap(), "share-orders");
806 assert_eq!(decoder.read_compact_string().unwrap(), "member-1");
807 assert_eq!(decoder.read_i32().unwrap(), 3);
808 assert_eq!(
809 decoder.read_compact_nullable_string().unwrap(),
810 Some("rack-a".to_owned())
811 );
812 assert_eq!(
813 decoder
814 .read_compact_array("topics", |decoder| decoder.read_compact_string())
815 .unwrap(),
816 Some(vec!["orders".to_owned()])
817 );
818 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
819 assert!(decoder.is_empty());
820 }
821
822 #[test]
823 fn decodes_share_group_heartbeat_v1_assignment() {
824 let mut bytes = Encoder::new();
825 bytes.write_i32(9);
826 bytes.write_i16(0);
827 bytes.write_compact_nullable_string(None).unwrap();
828 bytes
829 .write_compact_nullable_string(Some("member-1"))
830 .unwrap();
831 bytes.write_i32(4);
832 bytes.write_i32(2500);
833 bytes.write_i8(1);
834 bytes
835 .write_compact_array(
836 Some(&[ShareTopicPartitionsV1 {
837 topic_id: [7; 16],
838 partitions: vec![0, 2],
839 }]),
840 |encoder, assignment| assignment.encode(encoder),
841 )
842 .unwrap();
843 bytes.write_empty_tagged_fields();
844 bytes.write_empty_tagged_fields();
845
846 let encoded = bytes.into_bytes();
847 let mut decoder = Decoder::new(&encoded);
848 let response = ShareGroupHeartbeatResponseV1::decode_body(&mut decoder).unwrap();
849 assert_eq!(response.member_id.as_deref(), Some("member-1"));
850 assert_eq!(response.member_epoch, 4);
851 assert_eq!(
852 response.assignment.unwrap().topic_partitions[0].partitions,
853 vec![0, 2]
854 );
855 assert!(decoder.is_empty());
856 }
857
858 #[test]
859 fn encodes_share_fetch_v1_request_with_acknowledgements() {
860 let request = ShareFetchRequestV1 {
861 correlation_id: 22,
862 client_id: Some("kafrust".to_owned()),
863 group_id: Some("share-orders".to_owned()),
864 member_id: Some("member-1".to_owned()),
865 share_session_epoch: 2,
866 max_wait_ms: 500,
867 min_bytes: 1,
868 max_bytes: 1024,
869 max_records: 100,
870 batch_size: 10,
871 topics: vec![ShareFetchTopicV1 {
872 topic_id: [3; 16],
873 partitions: vec![ShareFetchPartitionV1 {
874 partition_index: 0,
875 acknowledgement_batches: vec![ShareAcknowledgementBatchV1 {
876 first_offset: 10,
877 last_offset: 12,
878 acknowledgement_types: vec![1, 2, 3],
879 }],
880 }],
881 }],
882 forgotten_topics: vec![ShareForgottenTopicV1 {
883 topic_id: [4; 16],
884 partitions: vec![1],
885 }],
886 };
887
888 let encoded = request.encode().unwrap();
889 assert_eq!(&encoded[0..4], &[0, 78, 0, 1]);
890 assert!(encoded.ends_with(&[0]));
891 let mut decoder = Decoder::new(&encoded[4..]);
892 assert_eq!(decoder.read_i32().unwrap(), 22);
893 assert_eq!(
894 decoder.read_nullable_string().unwrap(),
895 Some("kafrust".to_owned())
896 );
897 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
898 assert_eq!(
899 decoder.read_compact_nullable_string().unwrap(),
900 Some("share-orders".to_owned())
901 );
902 assert_eq!(
903 decoder.read_compact_nullable_string().unwrap(),
904 Some("member-1".to_owned())
905 );
906 assert_eq!(decoder.read_i32().unwrap(), 2);
907 assert_eq!(decoder.read_i32().unwrap(), 500);
908 assert_eq!(decoder.read_i32().unwrap(), 1);
909 assert_eq!(decoder.read_i32().unwrap(), 1024);
910 assert_eq!(decoder.read_i32().unwrap(), 100);
911 assert_eq!(decoder.read_i32().unwrap(), 10);
912 assert!(decoder
913 .read_compact_array("topics", |decoder| {
914 let topic_id = decoder.read_uuid()?;
915 let partitions = decoder.read_compact_array("partitions", |decoder| {
916 let partition_index = decoder.read_i32()?;
917 let batches = decoder
918 .read_compact_array("batches", ShareAcknowledgementBatchV1::decode)?;
919 decoder.read_tagged_fields()?;
920 Ok((partition_index, batches))
921 })?;
922 decoder.read_tagged_fields()?;
923 Ok((topic_id, partitions))
924 })
925 .unwrap()
926 .is_some());
927 }
928
929 #[test]
930 fn encodes_share_fetch_v2_request_with_record_limit_mode() {
931 let request = ShareFetchRequestV2 {
932 correlation_id: 23,
933 client_id: Some("kafrust".to_owned()),
934 group_id: Some("share-orders".to_owned()),
935 member_id: Some("member-1".to_owned()),
936 share_session_epoch: 3,
937 max_wait_ms: 500,
938 min_bytes: 1,
939 max_bytes: 1024,
940 max_records: 6,
941 batch_size: 2,
942 share_acquire_mode: 1,
943 is_renew_ack: false,
944 topics: Vec::new(),
945 forgotten_topics: Vec::new(),
946 };
947
948 let encoded = request.encode().unwrap();
949 assert_eq!(&encoded[0..4], &[0, 78, 0, 2]);
950 let mut decoder = Decoder::new(&encoded[4..]);
951 assert_eq!(decoder.read_i32().unwrap(), 23);
952 assert_eq!(
953 decoder.read_nullable_string().unwrap(),
954 Some("kafrust".to_owned())
955 );
956 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
957 assert_eq!(
958 decoder.read_compact_nullable_string().unwrap(),
959 Some("share-orders".to_owned())
960 );
961 assert_eq!(
962 decoder.read_compact_nullable_string().unwrap(),
963 Some("member-1".to_owned())
964 );
965 assert_eq!(decoder.read_i32().unwrap(), 3);
966 assert_eq!(decoder.read_i32().unwrap(), 500);
967 assert_eq!(decoder.read_i32().unwrap(), 1);
968 assert_eq!(decoder.read_i32().unwrap(), 1024);
969 assert_eq!(decoder.read_i32().unwrap(), 6);
970 assert_eq!(decoder.read_i32().unwrap(), 2);
971 assert_eq!(decoder.read_i8().unwrap(), 1);
972 assert!(!decoder.read_bool().unwrap());
973 assert!(decoder
974 .read_compact_array("topics", |_decoder| Ok::<(), Error>(()))
975 .unwrap()
976 .is_some());
977 assert!(decoder
978 .read_compact_array("forgotten topics", |_decoder| Ok::<(), Error>(()))
979 .unwrap()
980 .is_some());
981 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
982 assert!(decoder.is_empty());
983 }
984
985 #[test]
986 fn encodes_share_fetch_v2_renew_request_with_zero_fetch_limits() {
987 let request = ShareFetchRequestV2 {
988 correlation_id: 24,
989 client_id: Some("kafrust".to_owned()),
990 group_id: Some("share-orders".to_owned()),
991 member_id: Some("member-1".to_owned()),
992 share_session_epoch: 7,
993 max_wait_ms: 0,
994 min_bytes: 0,
995 max_bytes: 0,
996 max_records: 0,
997 batch_size: 0,
998 share_acquire_mode: 0,
999 is_renew_ack: true,
1000 topics: Vec::new(),
1001 forgotten_topics: Vec::new(),
1002 };
1003
1004 let encoded = request.encode().unwrap();
1005 let mut decoder = Decoder::new(&encoded[4..]);
1006 assert_eq!(decoder.read_i32().unwrap(), 24);
1007 assert_eq!(
1008 decoder.read_nullable_string().unwrap(),
1009 Some("kafrust".to_owned())
1010 );
1011 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
1012 assert_eq!(
1013 decoder.read_compact_nullable_string().unwrap(),
1014 Some("share-orders".to_owned())
1015 );
1016 assert_eq!(
1017 decoder.read_compact_nullable_string().unwrap(),
1018 Some("member-1".to_owned())
1019 );
1020 assert_eq!(decoder.read_i32().unwrap(), 7);
1021 assert_eq!(decoder.read_i32().unwrap(), 0);
1022 assert_eq!(decoder.read_i32().unwrap(), 0);
1023 assert_eq!(decoder.read_i32().unwrap(), 0);
1024 assert_eq!(decoder.read_i32().unwrap(), 0);
1025 assert_eq!(decoder.read_i32().unwrap(), 0);
1026 assert_eq!(decoder.read_i8().unwrap(), 0);
1027 assert!(decoder.read_bool().unwrap());
1028 assert!(decoder
1029 .read_compact_array("topics", |_decoder| Ok::<(), Error>(()))
1030 .unwrap()
1031 .is_some());
1032 assert!(decoder
1033 .read_compact_array("forgotten topics", |_decoder| Ok::<(), Error>(()))
1034 .unwrap()
1035 .is_some());
1036 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
1037 assert!(decoder.is_empty());
1038 }
1039
1040 #[test]
1041 fn decodes_share_fetch_v1_response_and_preserves_record_bytes() {
1042 let mut bytes = Encoder::new();
1043 bytes.write_i32(7);
1044 bytes.write_i16(0);
1045 bytes.write_compact_nullable_string(None).unwrap();
1046 bytes.write_i32(30_000);
1047 bytes
1048 .write_compact_array(
1049 Some(&[ShareFetchTopicResponseV1 {
1050 topic_id: [5; 16],
1051 partitions: vec![ShareFetchPartitionResponseV1 {
1052 partition_index: 0,
1053 error_code: 0,
1054 error_message: None,
1055 acknowledgement_error_code: 0,
1056 acknowledgement_error_message: None,
1057 current_leader: ShareLeaderIdAndEpochV1 {
1058 leader_id: 1,
1059 leader_epoch: 8,
1060 },
1061 records: Some(vec![1, 2, 3]),
1062 acquired_records: vec![ShareAcquiredRecordsV1 {
1063 first_offset: 10,
1064 last_offset: 12,
1065 delivery_count: 1,
1066 }],
1067 }],
1068 }]),
1069 |encoder, topic| {
1070 encoder.write_uuid(&topic.topic_id);
1071 encoder.write_compact_array(
1072 Some(&topic.partitions),
1073 |encoder, partition| {
1074 encoder.write_i32(partition.partition_index);
1075 encoder.write_i16(partition.error_code);
1076 encoder.write_compact_nullable_string(
1077 partition.error_message.as_deref(),
1078 )?;
1079 encoder.write_i16(partition.acknowledgement_error_code);
1080 encoder.write_compact_nullable_string(
1081 partition.acknowledgement_error_message.as_deref(),
1082 )?;
1083 encoder.write_i32(partition.current_leader.leader_id);
1084 encoder.write_i32(partition.current_leader.leader_epoch);
1085 encoder.write_empty_tagged_fields();
1086 encoder.write_compact_nullable_bytes(partition.records.as_deref())?;
1087 encoder.write_compact_array(
1088 Some(&partition.acquired_records),
1089 |encoder, acquired| {
1090 encoder.write_i64(acquired.first_offset);
1091 encoder.write_i64(acquired.last_offset);
1092 encoder.write_i16(acquired.delivery_count);
1093 encoder.write_empty_tagged_fields();
1094 Ok(())
1095 },
1096 )?;
1097 encoder.write_empty_tagged_fields();
1098 Ok(())
1099 },
1100 )?;
1101 encoder.write_empty_tagged_fields();
1102 Ok(())
1103 },
1104 )
1105 .unwrap();
1106 bytes
1107 .write_compact_array::<ShareNodeEndpointV1>(Some(&[]), |_encoder, _| Ok(()))
1108 .unwrap();
1109 bytes.write_empty_tagged_fields();
1110
1111 let encoded = bytes.into_bytes();
1112 let mut decoder = Decoder::new(&encoded);
1113 let response = ShareFetchResponseV1::decode_body(&mut decoder).unwrap();
1114 let partition = &response.responses[0].partitions[0];
1115 assert_eq!(partition.records.as_deref(), Some(&[1, 2, 3][..]));
1116 assert_eq!(partition.acquired_records[0].delivery_count, 1);
1117 assert!(decoder.is_empty());
1118 }
1119
1120 #[test]
1121 fn encodes_share_acknowledge_v1_request() {
1122 let request = ShareAcknowledgeRequestV1 {
1123 correlation_id: 33,
1124 client_id: None,
1125 group_id: Some("share-orders".to_owned()),
1126 member_id: Some("member-1".to_owned()),
1127 share_session_epoch: 3,
1128 topics: vec![ShareAcknowledgeTopicV1 {
1129 topic_id: [9; 16],
1130 partitions: vec![ShareAcknowledgePartitionV1 {
1131 partition_index: 2,
1132 acknowledgement_batches: vec![ShareAcknowledgementBatchV1 {
1133 first_offset: 20,
1134 last_offset: 20,
1135 acknowledgement_types: vec![1],
1136 }],
1137 }],
1138 }],
1139 };
1140
1141 let encoded = request.encode().unwrap();
1142 assert_eq!(&encoded[0..4], &[0, 79, 0, 1]);
1143 assert!(encoded.ends_with(&[0]));
1144 }
1145
1146 #[test]
1147 fn encodes_share_acknowledge_v2_renew_request() {
1148 let request = ShareAcknowledgeRequestV2 {
1149 correlation_id: 34,
1150 client_id: None,
1151 group_id: Some("share-orders".to_owned()),
1152 member_id: Some("member-1".to_owned()),
1153 share_session_epoch: 3,
1154 is_renew_ack: true,
1155 topics: vec![ShareAcknowledgeTopicV1 {
1156 topic_id: [9; 16],
1157 partitions: vec![ShareAcknowledgePartitionV1 {
1158 partition_index: 2,
1159 acknowledgement_batches: vec![ShareAcknowledgementBatchV1 {
1160 first_offset: 20,
1161 last_offset: 20,
1162 acknowledgement_types: vec![4],
1163 }],
1164 }],
1165 }],
1166 };
1167
1168 let encoded = request.encode().unwrap();
1169 assert_eq!(&encoded[0..4], &[0, 79, 0, 2]);
1170 let mut decoder = Decoder::new(&encoded[4..]);
1171 assert_eq!(decoder.read_i32().unwrap(), 34);
1172 assert_eq!(decoder.read_nullable_string().unwrap(), None);
1173 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
1174 assert_eq!(
1175 decoder.read_compact_nullable_string().unwrap(),
1176 Some("share-orders".to_owned())
1177 );
1178 assert_eq!(
1179 decoder.read_compact_nullable_string().unwrap(),
1180 Some("member-1".to_owned())
1181 );
1182 assert_eq!(decoder.read_i32().unwrap(), 3);
1183 assert!(decoder.read_bool().unwrap());
1184 let topics = decoder
1185 .read_compact_array("topics", |decoder| {
1186 let topic_id = decoder.read_uuid()?;
1187 let partitions = decoder
1188 .read_compact_array("partitions", |decoder| {
1189 let partition_index = decoder.read_i32()?;
1190 let batches = decoder
1191 .read_compact_array("batches", ShareAcknowledgementBatchV1::decode)?;
1192 decoder.read_tagged_fields()?;
1193 Ok((partition_index, batches.unwrap_or_default()))
1194 })?
1195 .unwrap_or_default();
1196 decoder.read_tagged_fields()?;
1197 Ok((topic_id, partitions))
1198 })
1199 .unwrap()
1200 .unwrap();
1201 assert_eq!(topics[0].1[0].1[0].acknowledgement_types, vec![4]);
1202 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
1203 assert!(decoder.is_empty());
1204 }
1205
1206 #[test]
1207 fn decodes_share_acknowledge_v1_response_with_leader_endpoint() {
1208 let mut bytes = Encoder::new();
1209 bytes.write_i32(4);
1210 bytes.write_i16(0);
1211 bytes.write_compact_nullable_string(None).unwrap();
1212 bytes
1213 .write_compact_array::<ShareAcknowledgeTopicResponseV1>(Some(&[]), |_encoder, _| Ok(()))
1214 .unwrap();
1215 bytes
1216 .write_compact_array(
1217 Some(&[ShareNodeEndpointV1 {
1218 node_id: 2,
1219 host: "broker".to_owned(),
1220 port: 9092,
1221 rack: None,
1222 }]),
1223 |encoder, endpoint| {
1224 encoder.write_i32(endpoint.node_id);
1225 encoder.write_compact_string(&endpoint.host)?;
1226 encoder.write_i32(endpoint.port);
1227 encoder.write_compact_nullable_string(endpoint.rack.as_deref())?;
1228 encoder.write_empty_tagged_fields();
1229 Ok(())
1230 },
1231 )
1232 .unwrap();
1233 bytes.write_empty_tagged_fields();
1234
1235 let encoded = bytes.into_bytes();
1236 let mut decoder = Decoder::new(&encoded);
1237 let response = ShareAcknowledgeResponseV1::decode_body(&mut decoder).unwrap();
1238 assert_eq!(response.node_endpoints[0].host, "broker");
1239 assert_eq!(response.node_endpoints[0].port, 9092);
1240 assert!(decoder.is_empty());
1241 }
1242
1243 #[test]
1244 fn decodes_share_acknowledge_v2_response_with_lock_timeout() {
1245 let mut bytes = Encoder::new();
1246 bytes.write_i32(4);
1247 bytes.write_i16(0);
1248 bytes.write_compact_nullable_string(None).unwrap();
1249 bytes.write_i32(45_000);
1250 bytes
1251 .write_compact_array::<ShareAcknowledgeTopicResponseV1>(Some(&[]), |_encoder, _| Ok(()))
1252 .unwrap();
1253 bytes
1254 .write_compact_array::<ShareNodeEndpointV1>(Some(&[]), |_encoder, _| Ok(()))
1255 .unwrap();
1256 bytes.write_empty_tagged_fields();
1257
1258 let encoded = bytes.into_bytes();
1259 let mut decoder = Decoder::new(&encoded);
1260 let response = ShareAcknowledgeResponseV2::decode_body(&mut decoder).unwrap();
1261 assert_eq!(response.acquisition_lock_timeout_ms, 45_000);
1262 assert!(decoder.is_empty());
1263 }
1264}